PyODPS是MaxCompute的Python SDK,提供对MaxCompute对象的基本操作及DataFrame框架,支持在DataWorks或本地环境中使用Python进行大数据处理和分析。本文介绍如何在本地环境中安装PyODPS、连接MaxCompute、执行SQL查询及使用DataFrame API。
前提条件
在开始之前,请确保具备以下条件:
已安装Python 3.6或以上版本。
已创建阿里云账号,并拥有AccessKey(用于连接MaxCompute项目)。建议使用RAM用户的AccessKey。
安装PyODPS
安装PyODPS
Linux/macOS:
pip3 install pyodpsWindows:
pip install pyodps
如果安装时出现numpy或者pyarrow等依赖包安装错误,通常显示为C代码编译错误,可能是pip或setuptools版本过旧导致,请先升级pip和setuptools,然后重试:
# Linux/macOS pip3 install -U pip setuptools # Windows pip install -U pip setuptools验证安装
Linux/macOS:
python3 -c "from odps import ODPS"Windows:
python -c "from odps import ODPS"无报错信息即表示安装成功。如果出现错误,请参见问题排查。
设置环境变量
AccessKey是访问MaxCompute项目的身份凭证,请将其配置为环境变量,避免在代码中硬编码。
登录RAM控制台,复制AccessKey ID和AccessKey Secret,然后根据操作系统配置。macOS(zsh)配置方式如下:
macOS(zsh):
vim ~/.zshrc添加以下内容,将占位符替换为实际值
export ALIBABA_CLOUD_ACCESS_KEY_ID=<your-AccessKey-ID> export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<your-AccessKey-secret>使配置生效
source ~/.zshrc确认环境变量已设置
echo $ALIBABA_CLOUD_ACCESS_KEY_IDecho $ALIBABA_CLOUD_ACCESS_KEY_SECRET每条命令会输出对应的值。如果输出为空,请重复步骤2和3。
关于在Linux、macOS和Windows系统上设置环境变量的更多详细信息,请参见在环境变量中设置阿里云AccessKey。
初始化ODPS入口
ODPS类是所有PyODPS操作的入口。在调用任何其他API之前,需要使用凭证和项目信息初始化。
运行代码前,请确保已设置ALIBABA_CLOUD_ACCESS_KEY_ID和ALIBABA_CLOUD_ACCESS_KEY_SECRET环境变量。如果未设置环境变量,也可以显式指定AccessKey,但是在源代码中硬编码凭证存在安全风险,所以不推荐该方式。
如需获取AccessKey ID和AccessKey Secret,请单击链接获取ACCESS_ID及SECRET_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,请填写 |
完成上述配置后,即可在本地使用 PyODPS。更多操作说明请参见基本操作概述。
数据处理
SQL
MaxCompute支持两种SQL执行模式:传统模式和查询加速MaxQA模式。两种模式均支持DDL和DML语句。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)
查询加速MaxQA是MaxCompute的查询加速功能,使用独立的资源池加速中小规模数据查询。从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操作。
本地创建
test-pyodps-local.py文件。添加以下代码:
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()运行脚本
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代码编译错误。
原因: pip或setuptools版本过旧。
解决方案: 升级pip和setuptools后重新安装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"