在PAI子产品(DLC或DSW)中,您可以使用MaxCompute的PyODPS SDK或PAI自研的paiio模块读写MaxCompute数据。根据是否使用TensorFlow、是否需要即时I/O等场景,选择合适的方式。
功能介绍
PyODPS是MaxCompute的Python SDK。您可以用它上传下载文件、创建表、运行ODPS SQL查询等。详情请参见PyODPS概述。
PAI自研的paiio模块提供三种接口,方便读写MaxCompute表:
接口
区别
功能描述
TableRecordDataset
依赖TensorFlow。推荐在1.2及以上版本中使用Dataset接口(详情请参见Dataset)构建数据流,替代线程和队列接口。
读取MaxCompute表数据。
TableReader
不依赖TensorFlow,基于MaxCompute SDK实现,可直接访问MaxCompute表并即时获取I/O结果。
读取MaxCompute表数据。
TableWriter
不依赖TensorFlow,基于MaxCompute SDK实现,可直接写入MaxCompute表。
往MaxCompute表中写入数据。
前提条件
已安装Python 3.6及以上版本。Python 2.7及以下版本不支持。
已配置环境变量。具体操作,请参见在Linux、macOS和Windows系统配置环境变量。
已开通MaxCompute并创建项目。详情请参见开通MaxCompute和创建MaxCompute项目。
使用限制
paiio模块不支持自定义镜像。仅当选择TensorFlow 1.12、1.15或2.0镜像时可用。
PyODPS
您可以使用PyODPS读写MaxCompute数据。
执行如下命令安装PyODPS。
pip install pyodps执行如下命令检查安装是否成功。无返回值和报错即表示安装成功。
python -c "from odps import ODPS"如果使用的Python不是系统默认版本,安装pip后可执行如下命令切换Python版本。
/home/tops/bin/python3.7 -m pip install setuptools>=3.0 #/home/tops/bin/python3.7为安装的python路径通过PyODPS读写MaxCompute数据,示例代码如下。
import numpy as np import pandas as pd import os from odps import ODPS from odps.df import DataFrame # 建立链接。 o = ODPS( os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID'), os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET'), project='your-default-project', endpoint='your-end-point', ) # 直接指定字段名和类型,创建非分区表my_new_table。 table = o.create_table('my_new_table', 'num bigint, id string', if_not_exists=True) # 向my_new_table中插入数据。 records = [[111, 'aaa'], [222, 'bbb'], [333, 'ccc'], [444, '中文']] o.write_table(table, records) # 读取数据。 sql = ''' SELECT * FROM your-default-project.<table> LIMIT 100 ; ''' query_job = o.execute_sql(sql) result = query_job.open_reader(tunnel=True) df = result.to_pandas(n_process=1) # n_process根据机器配置设置,大于1时开启多线程加速。
paiio
准备工作:配置账户信息
使用paiio前,需配置MaxCompute账户的AccessKey信息。您可以将配置文件放在挂载的文件系统中,再通过环境变量引用。
编写配置文件,内容如下。
access_id=xxxx access_key=xxxx end_point=http://xxxx参数
描述
access_id
阿里云账号的AccessKey ID。
access_key
阿里云账号的AccessKey Secret。
end_point
MaxCompute的Endpoint,例如华东2(上海)配置为
http://service.cn-shanghai.maxcompute.aliyun.com/api。详情请参见Endpoint。在代码中指定配置文件路径。
os.environ['ODPS_CONFIG_FILE_PATH'] = '<your MaxCompute config file path>'其中<your MaxCompute config file path>表示配置文件的路径。
TableRecordDataset使用说明
接口说明
TensorFlow社区推荐在1.2及以上版本中使用Dataset接口(详情请参见Dataset)替代原有的线程和队列接口构建数据流。通过多个Dataset接口的组合变换生成计算数据,可以简化数据输入部分的代码。
接口定义(Python)
class TableRecordDataset(Dataset): def __init__(self, filenames, record_defaults, selected_cols=None, excluded_cols=None, slice_id=0, slice_count=1, num_threads=0, capacity=0):参数
参数
是否必选
类型
默认值
描述
filenames
是
STRING
无
待读取的表名列表,多张表的Schema必须一致。表名格式为
odps://${your_projectname}/tables/${table_name}/${pt_1}/${pt_2}/...。record_defaults
是
LIST或TUPLE
无
读取列的数据类型转换及列为空时的默认值。列数或数据类型不匹配时会抛出异常。
支持FLOAT32、FLOAT64、INT32、INT64、BOOL及STRING。INT64默认值请参见
np.array(0, np.int64)。selected_cols
否
STRING
None
要读取的列,用半角逗号分隔。默认None表示读取所有列。不能与excluded_cols同时使用。
excluded_cols
否
STRING
None
要排除的列,用半角逗号分隔。默认None表示读取所有列。不能与selected_cols同时使用。
slice_id
否
INT
0
分布式读取场景下当前分片的编号,从0开始。系统根据slice_count将表均分为多个分片,读取slice_id对应的分片。
slice_id为0且slice_count为1时读取整张表;slice_count大于1时读取第0个分片。
slice_count
否
INT
1
分布式读取场景下的总分片数,通常为Worker数量。默认1表示不分片,读取整张表。
num_threads
否
INT
0
预取数据时每个Reader启用的线程数,独立于计算线程,取值1~64。取num_threads为0时,系统自动将其设为计算线程池线程数的1/4。
说明I/O对模型整体计算的影响因模型而异,提高预取线程数不一定能提升训练速度。
capacity
否
INT
0
读取表的总预取行数。若num_threads大于1,每个线程预取capacity/num_threads行(向上取整)。若capacity为0,Reader根据前256行的平均行大小自动配置总预取量,使每个线程预取数据约64 MB。
说明如果MaxCompute表字段为DOUBLE类型,则TensorFlow中需要使用np.float64格式与其对应。
返回值
返回Dataset对象,可作为Pipeline的输入。
使用示例
假设myproject项目中有一张test表,内容如下。
itemid(BIGINT) | name(STRING) | price(DOUBLE) | virtual(BOOL) |
25 | "Apple" | 5.0 | False |
38 | "Pear" | 4.5 | False |
17 | "Watermelon" | 2.2 | False |
以下代码使用TableRecordDataset读取test表的itemid和price列。
import os
import tensorflow as tf
import paiio
# 指定配置文件路径,请替换为实际路径。
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# 定义要读取的表,可填多个。请替换为实际表名和MaxCompute项目名称。
table = ["odps://${your_projectname}/tables/${table_name}"]
# 定义TableRecordDataset,读取itemid和price列。
dataset = paiio.data.TableRecordDataset(table,
record_defaults=[0, 0.0],
selected_cols="itemid,price",
num_threads=1,
capacity=10)
# 设置epoch=2,batch_size=3,prefetch=100。
dataset = dataset.repeat(2).batch(3).prefetch(100)
ids, prices = tf.compat.v1.data.make_one_shot_iterator(dataset).get_next()
with tf.compat.v1.Session() as sess:
sess.run(tf.compat.v1.global_variables_initializer())
sess.run(tf.compat.v1.local_variables_initializer())
try:
while True:
batch_ids, batch_prices = sess.run([ids, prices])
print("batch_ids:", batch_ids)
print("batch_prices:", batch_prices)
except tf.errors.OutOfRangeError:
print("End of dataset")TableReader使用说明
接口说明
TableReader基于MaxCompute SDK实现,不依赖TensorFlow,可直接访问MaxCompute表并即时获取I/O结果。
创建Reader并打开表
接口定义
reader = paiio.python_io.TableReader(table, selected_cols="", excluded_cols="", slice_id=0, slice_count=1):参数
返回值
Reader对象。
参数 | 是否必选 | 类型 | 默认值 | 描述 |
table | 是 | STRING | 无 | 要打开的MaxCompute表名,格式为 |
selected_cols | 否 | STRING | 空字符串("") | 要读取的列,用英文逗号分隔。默认空字符串("")表示读取所有列。不能与excluded_cols同时使用。 |
excluded_cols | 否 | STRING | 空字符串("") | 要排除的列,用英文逗号分隔。默认空字符串("")表示读取所有列。不能与selected_cols同时使用。 |
slice_id | 否 | INT | 0 | 分布式读取场景下当前分片的编号,取值范围[0, slice_count-1]。系统根据slice_count将表均分为多个分片,读取slice_id对应的分片。默认0表示不分片,读取所有行。 |
slice_count | 否 | INT | 1 | 分布式读取场景下的总分片数,通常为Worker数量。 |
读取记录
接口定义
reader.read(num_records=1)参数
num_records表示顺序读取的行数,默认1。若超出未读行数,返回剩余所有行;未读到记录则抛出OutOfRange异常(paiio.python_io.OutOfRangeException)。
返回值
返回numpy ndarray(或recarray),每个元素为一行数据组成的TUPLE。
定位到相应行
接口定义
reader.seek(offset=0)参数
offset表示定位到的行号,从0开始,下一个read从该行开始。若配置了slice_id和slice_count,则按分片内相对位置定位。offset超出总行数,或已读到表尾后继续seek,都会抛出OutOfRange异常(paiio.python_io.OutOfRangeException)。
读取一个batch时,若剩余行数不足batch_size,read返回剩余行且不抛异常;此时继续seek会抛异常。
返回值
无返回值。操作出错时抛出异常。
获取表的总记录数
接口定义
reader.get_row_count()参数
无
返回值
返回表的行数。若配置了slice_id和slice_count,返回当前分片大小。
获取表的Schema
接口定义
reader.get_schema()参数
无
返回值
返回1D structured ndarray,每个元素对应一列的Schema,包括以下三项。
参数 | 描述 |
colname | 列名 |
typestr | MaxCompute数据类型名称。 |
pytype | typestr对应的Python类型。 |
typestr和pytype的对应关系如下表所示。
typestr | pytype |
BIGINT | INT |
DOUBLE | FLOAT |
BOOLEAN | BOOL |
STRING | OBJECT |
DATETIME | INT |
MAP 说明 PAI-TensorFlow不支持对MAP类型数据进行操作。 | OBJECT |
关闭表
接口定义
reader.close()参数
无
返回值
无返回值。操作出错时抛出异常。
使用示例
假设myproject项目中有一张test表,内容如下。
uid(BIGINT) | name(STRING) | price(DOUBLE) | virtual(BOOL) |
25 | "Apple" | 5.0 | False |
38 | "Pear" | 4.5 | False |
17 | "Watermelon" | 2.2 | False |
以下代码使用TableReader读取test表的uid、name和price列。
import os
import paiio
# 指定配置文件路径,请替换为实际路径。
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# 打开表并返回reader对象,请替换为实际表名和MaxCompute项目名称。
reader = paiio.python_io.TableReader("odps://myproject/tables/test", selected_cols="uid,name,price")
# 获取表的总行数。
total_records_num = reader.get_row_count() # return 3
batch_size = 2
# 读表,返回值将是一个recarray数组,形式为[(uid, name, price)*2]。
records = reader.read(batch_size) # 返回[(25, "Apple", 5.0), (38, "Pear", 4.5)]
records = reader.read(batch_size) # 返回[(17, "Watermelon", 2.2)]
# 继续读取将抛出OutOfRange异常。
# Close the reader.
reader.close()TableWriter使用说明
TableWriter基于MaxCompute SDK实现,不依赖TensorFlow,可直接写入MaxCompute表。
接口说明
创建Writer并打开表
接口定义
writer = paiio.python_io.TableWriter(table, slice_id=0)说明该接口不会清空原表数据,以追加方式写入。
新写入的数据需关闭表后才能读取。
参数
参数
是否必选
类型
默认值
描述
table
是
STRING
无
要打开的MaxCompute表名,格式为
odps://${your_projectname}/tables/${table_name}/${pt_1}/${pt_2}/...。slice_id
否
INT
0
分布式场景下将数据写入不同分片,避免写冲突。单机场景使用默认值0即可。多机场景下,多个Worker(包括PS)使用同一个slice_id会导致写入失败。
返回值
返回Writer对象。
写入记录
接口定义
writer.write(values, indices)参数
参数
是否必选
类型
默认值
描述
values
是
STRING
无
要写入的数据,支持单行或多行:
单行数据:传入由标量组成的TUPLE、LIST或1D-ndarray。传入LIST或ndarray时,各列数据类型需一致。
N行数据(N>=1):传入LIST或1D-ndarray,每个元素对应一行数据(用TUPLE或LIST表示,也可通过Structure存放于ndarray中)。
indices
是
INT
无
指定数据写入的列,传入由INT组成的TUPLE、LIST或1D-ndarray。indices中每个数i对应表中第i列,列号从0开始。
返回值
无返回值。写入出错时抛出异常。
关闭表
接口定义
writer.close()说明使用with语句时,无需显式调用close()。
参数
无
返回值
无返回值。操作出错时抛出异常。
示例
通过with语句使用TableWriter的代码如下。
with paiio.python_io.TableWriter(table) as writer: # 准备待写入的数据。 writer.write(values, indices) # 离开with块后,表会自动关闭。
使用示例
import paiio
import os
# 指定配置文件路径,请替换为实际路径。
os.environ['ODPS_CONFIG_FILE_PATH'] = "/mnt/data/odps_config.ini"
# 准备待写入的数据。
values = [(25, "Apple", 5.0, False),
(38, "Pear", 4.5, False),
(17, "Watermelon", 2.2, False)]
# 打开表并返回writer对象,请替换为实际表名和MaxCompute项目名称。
writer = paiio.python_io.TableWriter("odps://project/tables/test")
# 将数据写入表的第0-3列。
records = writer.write(values, indices=[0, 1, 2, 3])
# 关闭writer,确保数据写入完成。
writer.close()