在本地环境使用PyODPS

更新时间:
复制 MD 格式

PyODPSMaxComputePython SDK,提供对MaxCompute对象的基本操作及DataFrame框架,支持在DataWorks或本地环境中使用Python进行大数据处理和分析。本文介绍如何在本地环境中安装PyODPS、连接MaxCompute、执行SQL查询及使用DataFrame API。

前提条件

在开始之前,请确保具备以下条件:

  • 已安装Python 3.6或以上版本。

  • 已创建阿里云账号,并拥有AccessKey(用于连接MaxCompute项目)。建议使用RAM用户的AccessKey。

安装PyODPS

  1. 安装PyODPS

    • Linux/macOS:

      pip3 install pyodps
    • Windows:

      pip install pyodps

    如果安装时出现numpy或者pyarrow等依赖包安装错误,通常显示为C代码编译错误,可能是pipsetuptools版本过旧导致,请先升级pipsetuptools,然后重试:

    # Linux/macOS
    pip3 install -U pip setuptools
    
    # Windows
    pip install -U pip setuptools
  2. 验证安装

    • Linux/macOS:

      python3 -c "from odps import ODPS"
    • Windows:

      python -c "from odps import ODPS"

      无报错信息即表示安装成功。如果出现错误,请参见问题排查。

设置环境变量

AccessKey是访问MaxCompute项目的身份凭证,请将其配置为环境变量,避免在代码中硬编码。

登录RAM控制台,复制AccessKey IDAccessKey Secret,然后根据操作系统配置。macOS(zsh)配置方式如下:

  1. macOS(zsh):

    vim ~/.zshrc
  2. 添加以下内容,将占位符替换为实际值

    export ALIBABA_CLOUD_ACCESS_KEY_ID=<your-AccessKey-ID>
    export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<your-AccessKey-secret>
  3. 使配置生效

    source ~/.zshrc
  4. 确认环境变量已设置

    echo $ALIBABA_CLOUD_ACCESS_KEY_IDecho $ALIBABA_CLOUD_ACCESS_KEY_SECRET

    每条命令会输出对应的值。如果输出为空,请重复步骤23。

关于在Linux、macOSWindows系统上设置环境变量的更多详细信息,请参见在环境变量中设置阿里云AccessKey

初始化ODPS入口

ODPS类是所有PyODPS操作的入口。在调用任何其他API之前,需要使用凭证和项目信息初始化。

运行代码前,请确保已设置ALIBABA_CLOUD_ACCESS_KEY_IDALIBABA_CLOUD_ACCESS_KEY_SECRET环境变量。如果未设置环境变量,也可以显式指定AccessKey,但是在源代码中硬编码凭证存在安全风险,所以不推荐该方式。

如需获取AccessKey IDAccessKey Secret,请单击链接获取ACCESS_IDSECRET_ACCESS_KEY

手动定义ODPS入口

手动定义ODPS入口,代码示例如下:

  • 推荐写法(环境变量)

    import os
    from odps import ODPS
    
    o = ODPS(
        access_id=os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
        secret_access_key=os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
        project='your-default-project',
        endpoint='your-end-point',
    )
  • 不推荐写法(显式指定,仅用于测试)

    import os
    from odps import ODPS
    
    o = ODPS(
        access_id='your-aliyun-access-key-id',
        secret_access_key='your-aliyun-access-key-secret',
        project='your-default-project',
        endpoint='your-end-point',
    )
    

使用 STS 安全凭证访问 ODPS

若使用 STS 安全凭证访问 ODPS,可以使用下面的语句初始化 ODPS 入口对象:

import os
from odps import ODPS
from odps.accounts import StsAccount
# 确保 ALIBABA_CLOUD_ACCESS_KEY_ID 环境变量设置为 Access Key ID,
# ALIBABA_CLOUD_ACCESS_KEY_SECRET 环境变量设置为 Access Key Secret,
# ALIBABA_CLOUD_STS_TOKEN 环境变量设置为 STS Token,
# 不建议直接使用 Access Key ID / Access Key Secret 字符串
account = StsAccount(
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
    os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
    os.getenv('ALIBABA_CLOUD_STS_TOKEN'),
)
o = ODPS(
    account=account,
    project='**your-default-project**',
    endpoint='**your-end-point**',
)

参数

填写说明

your-default-project

填写MaxCompute项目名称,登录MaxCompute控制台获取项目列表。

your-end-point

填写所在地域和所在网络环境下的Endpoint

例如杭州地域的公网Endpoint,请填写https://service.cn-hangzhou.maxcompute.aliyun.com/api

完成上述配置后,即可在本地使用 PyODPS。更多操作说明请参见基本操作概述

数据处理

SQL

MaxCompute支持两种SQL执行模式:传统模式查询加速MaxQA模式。两种模式均支持DDLDML语句。MCQA会缓存作业结果,重复执行时返回缓存结果以加速查询。计费规则请参见计算费用(按量付费)

execute_sql()run_sql()并非支持所有SQL语句类型。对于非DDL和非DML语句,请使用对应的方法直接调用,例如,CREATE TABLE语句请使用create_table()方法,API命令请使用run_xflow()execute_xflow()方法。

传统模式

使用execute_sql()执行SQL语句,使用open_reader()读取结果:

# 使用专用方法创建表
o.create_table('my_t', 'num bigint, id string', if_not_exists=True)

# 执行SELECT查询
result = o.execute_sql('SELECT * FROM pyodps_iris LIMIT 3')
with result.open_reader() as reader:
    for record in reader:
        print(record)

加速查询模式(MaxQA)

查询加速MaxQAMaxCompute的查询加速功能,使用独立的资源池加速中小规模数据查询。从PyODPS v0.11.4.1开始,可以使用execute_sql_interactive通过MaxQA执行SQL。

o.execute_sql_interactive('SELECT * FROM dual', fallback='all')

如果MCQA无法执行该SQL,系统会自动回退到传统模式。如需禁止回退,将fallback参数设置为False。也可以指定回退策略,以逗号分隔的字符串形式传入:

回退策略

说明

generic

发生未知错误时回退到传统模式。

noresource

资源不足时回退到传统模式。

upgrading

系统升级期间回退到传统模式。

timeout

发生超时时回退到传统模式。

unsupported

遇到MCQA不支持的场景时回退到传统模式。

默认策略为unsupported,upgrading,noresource,timeout。如需确保始终启用回退,设置为all(即generic,unsupported,upgrading,noresource,timeout的组合)。

# 仅在资源不足和不支持的场景下回退
o.execute_sql_interactive('SELECT * FROM dual', fallback="noresource,unsupported")

更多PyODPS节点的SQL相关操作详情请参见SQL

DataFrame

PyODPS提供DataFrame API用于数据处理。

pandas不同,PyODPS DataFrame是惰性执行的,即只有在调用execute()persist()等立即执行方法时才会触发计算。

from odps.df import DataFrame

# 将pyodps_iris表加载为DataFrame
iris = DataFrame(o.get_table('pyodps_iris'))

# 过滤行并打印结果——.execute()触发实际计算
for record in iris[iris.sepalwidth < 3].execute():
    print(record)

默认情况下,本地环境的PyODPS节点运行过程不会打印Logview等详细过程。如需开启详细输出:

from odps import options
options.verbose = True

设置运行参数hints

通过向SQL调用传入hints字典,可以设置单次查询的运行参数:

o.execute_sql('SELECT * FROM pyodps_iris', hints={'odps.sql.mapper.split.size': 16})

如需对会话中的所有SQL调用生效,可以全局配置hints:

from odps import options

# 后续所有execute_sql()调用将自动包含这些hints
options.sql.settings = {'odps.sql.mapper.split.size': 16}
o.execute_sql('SELECT * FROM pyodps_iris')

完整示例

以下示例覆盖完整流程:连接MaxCompute、创建表、写入数据、读取数据、查询pyodps_iris表并清理资源。

本示例最后调用了table.drop(),会永久删除表my_new_table如果需要保留该表,请勿使用已有的表名或执行drop操作。

  1. 本地创建test-pyodps-local.py文件。

  2. 添加以下代码:

    import os
    from odps import ODPS
    
    o = ODPS(
        access_id=os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'),
        secret_access_key=os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'),
        project='your-default-project',
        endpoint='your-end-point',
    )
    
    table = o.create_table('my_new_table', 'num bigint, id string', if_not_exists=True)
    
    records = [[111, 'aaa'],
               [222, 'bbb'],
               [333, 'ccc'],
               [444, '中文']]
    o.write_table(table, records)
    
    # 读取记录
    for record in o.read_table(table):
        print(record[0], record[1])
    
     # 使用SQL查询my_new_table表
    result = o.execute_sql('SELECT * FROM my_new_table;')
    
    
    print('Read data from the pyodps_iris table using open_reader:')
    
    # 打印查询结果
    with result.open_reader() as reader:
        for record in reader:
            print(record[0], record[1])
    
    # 删除表以释放资源
    table.drop()
  3. 运行脚本

    python test-pyodps-local.py

    预期运行结果:

    111 aaa
    222 bbb
    333 ccc
    444 中文
    Read data from the pyodps_iris table using open_reader:
    111 aaa
    222 bbb
    333 ccc
    444 中文

问题排查

依赖包安装错误

  • 现象: 安装numpy、pyarrow等依赖包时出现C代码编译错误。

  • 原因: pipsetuptools版本过旧。

  • 解决方案: 升级pipsetuptools后重新安装PyODPS。

    # Linux/macOS
    pip3 install -U pip setuptools
    
    # Windows
    pip install -U pip setuptools

Python版本下pip版本冲突

  • 现象: PyODPS安装到了错误的Python版本上。

  • 原因: 系统中存在多个Python版本,pip3指向的Python版本与预期不符。

  • 解决方案: 通过指定Python可执行文件来调用pip。

    /home/tops/bin/python3.7 -m pip install pyodps

    /home/tops/bin/python3.7替换为目标Python安装的路径。

urllib3 OpenSSL版本错误

  • 现象: 安装失败,提示urllib3 v2.0 only supports OpenSSL 1.1.1+

  • 原因: Python安装使用了较旧版本的OpenSSL,与urllib3 v2.0不兼容。

  • 解决方案: 先安装较低版本的urllib3,再安装PyODPS。

    # Linux/macOS
    pip3 install "urllib3<2.0"# Windows
    pip install "urllib3<2.0"