- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 提供了内置的 TFRecordIO 连接器,让 Python 管道可以像读取普通文本一样读取 TensorFlow 生态中广泛使用的 TFRecord 二进制格式。本文以仓库中 code-explanation 提示模板 06_io_tfrecord.md 讲解的代码为骨架,完整还原一个"命令行传参 → 读取 TFRecord → 反序列化 → 日志输出"的实战管道,并结合 tfrecordio.py 源码与 tfrecordio_test.py 测试,深入剖析ReadFromTFRecord的参数语义、TFRecord 的字节级编码格式与 CRC 校验机制。读完本文,你将能独立编写、运行并调试 Beam Python 的 TFRecord 读写管道。
一、TFRecord 是什么,为什么 Beam 需要专门的连接器
TFRecord 是 TensorFlow 推荐的一种二进制记录格式,常用于把大规模训练样本(图片、特征向量等)打包成便于顺序读取的文件。它本身不关心记录内部的字段含义——每个记录本质上就是一段字节,具体语义由写入方自行定义(例如用tf.train.Example协议消息编码,或用 Pythonpickle序列化)。
这种"一段字节 + 强校验"的格式特性决定了它不能像文本文件那样按行切分读取:每条记录带有长度与 CRC 校验信息,读取时必须逐条解析、逐条校验。因此 Apache Beam 在 Python SDK 中提供了专用的 TFRecordIO 连接器,即apache_beam.io.tfrecordio模块,对外暴露三个核心变换:
ReadFromTFRecord:从文件通配符(glob)读取 TFRecord 文件;ReadAllFromTFRecord:从元素为文件路径的PCollection读取,适合运行时才知道文件列表的场景;WriteToTFRecord:把PCollection写入 TFRecord 格式。
这三个类均在 tfrecordio.py 的__all__中导出。
二、完整可运行示例:从命令行参数读取 TFRecord
下面的完整版本取自配套的 code-generation 模板 06_io_tfrecord.md,它在原解释文档的基础上补齐了 import、main入口与日志级别设置,可以直接保存运行:
import logging import apache_beam as beam from apache_beam import Map from apache_beam.io.tfrecordio import ReadFromTFRecord from apache_beam.options.pipeline_options import PipelineOptions class TFRecordOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument( "--file_pattern", help="A file glob pattern to read TFRecords from." ) def map_from_bytes(element): """ Deserializes the input bytes using pickle library and returns the reconstructed object. """ # third party libraries import pickle return pickle.loads(element) def run(): options = TFRecordOptions() with beam.Pipeline(options=options) as p: output = ( p | "Read from TFRecord" >> ReadFromTFRecord( file_pattern=options.file_pattern ) | "Map from bytes" >> Map(map_from_bytes) | "Log Data" >> Map(logging.info) ) if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) run()运行方式(以 Direct Runner 本地执行为例):
python read_tfrecord.py --file_pattern="gs://your-bucket/data/train-*.tfrecord" # 或者本地文件 python read_tfrecord.py --file_pattern="./data/part-*.tfrecord"管道由三段组成:ReadFromTFRecord读出字节 →Map(map_from_bytes)用pickle.loads还原对象 →Map(logging.info)把每条对象打到控制台。
三、逐段拆解:管道选项、读取与反序列化
3.1 用 PipelineOptions 子类解析命令行参数
class TFRecordOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument( "--file_pattern", help="A file glob pattern of TFRecord files" ) options = TFRecordOptions()PipelineOptions是 Beam Python SDK 的标准命令行参数机制:继承它并覆写_add_argparse_args类方法,即可往 argparse 解析器中注册自定义参数。框架会自动完成参数解析,并把解析结果作为PipelineOptions的属性暴露。这里注册的--file_pattern被存为options.file_pattern,随后传给ReadFromTFRecord,这样文件路径就不需要硬编码在管道代码里,便于在不同环境(本地 / Dataflow / 其他 Runner)间复用同一份代码。
3.2 ReadFromTFRecord:按 glob 读取全部匹配文件
p | "Read from TFRecord" >> ReadFromTFRecord( file_pattern=options.file_pattern )ReadFromTFRecord接收一个文件 glob 通配模式(如data/*.tfrecord、train-?????.tfrecord),读取所有匹配文件并把每条记录作为一个元素输出。默认情况下,输出的每个元素是bytes类型——因为 TFRecord 变换默认使用coders.BytesCoder()。
3.3 Map(map_from_bytes):把字节还原成 Python 对象
def map_from_bytes(element): return pickle.loads(element)ReadFromTFRecord产出的是原始字节,业务层需要自行把字节解码成对象。示例中写入方用pickle.dumps序列化了对象,因此读取端用pickle.loads还原。这是 TFRecord 最典型的用法:Beam 负责记录边界与校验,语义解码交给用户自定义的Map/DoFn。
3.4 Map(logging.info):观察输出
| "Log Data" >> Map(logging.info)Map(logging.info)把每条元素直接作为参数传给logging.info。配合main里logging.getLogger().setLevel(logging.INFO),运行管道即可在控制台逐条看到反序列化后的对象内容,方便快速验证读取链路是否正常。
四、ReadFromTFRecord 参数详解(以源码为准)
ReadFromTFRecord的完整签名位于 tfrecordio.py:
class ReadFromTFRecord(PTransform): def __init__( self, file_pattern, coder=coders.BytesCoder(), compression_type=CompressionTypes.AUTO, validate=True):| 参数 | 类型 | 默认值 | 语义 |
|---|---|---|---|
file_pattern | str | 必填 | 要读取的 TFRecord 文件的 glob 通配模式,如data/*.tfrecord |
coder | Coder | coders.BytesCoder() | 用于解码每条记录的编码器;默认按原始字节输出 |
compression_type | str | CompressionTypes.AUTO | 压缩类型,AUTO表示按文件扩展名自动探测 |
validate | bool | True | 管道构建阶段即校验文件是否存在,提前暴露路径错误 |
几个值得注意的实现细节:
coder默认是BytesCoder,定义于 coders.py。如果 TFRecord 文件里存的本来就是字符串,可以换成coders.StrUtf8Coder()直接得到str;如果存的是tf.train.Example序列化字节,通常保持BytesCoder再由下游tf.train.Example.FromString解析。compression_type的可选值定义在 filesystem.py:AUTO、BZIP2、DEFLATE、ZSTD、GZIP、LZMA、UNCOMPRESSED。AUTO模式下,Beam 依据文件扩展名(如.gz、.bz2)自动判断压缩方式,未知扩展名按未压缩处理。validate=True会在管道创建阶段检查文件是否存在;当文件由上游任务动态生成、创建时尚未落盘时,应改为False避免校验失败。
从 tfrecordio_test.py 的TestReadFromTFRecord可以看到参数组合的标准用法,例如用BytesCoder配合CompressionTypes.GZIP、CompressionTypes.DEFLATE、CompressionTypes.AUTO分别读取压缩与非压缩文件,并断言输出为[b'foo', b'bar']之类的原始字节列表。
五、底层原理:TFRecord 的字节格式与 CRC 校验
为什么读取时不能简单按行 split?看_TFRecordUtil(tfrecordio.py)就明白了。TFRecord 文件中每条记录按固定布局写入(LittleEndian 字节序):
┌────────────────────┬──────────────────────┬──────────────┬──────────────────┐ │ 8 字节 长度 (uint64) │ 4 字节 长度 CRC32C │ 数据 (N 字节) │ 4 字节 数据 CRC32C │ └────────────────────┴──────────────────────┴──────────────┴──────────────────┘- 写入时
write_record先用struct.pack('<Q', len(value))写 8 字节长度,再分别对长度字段和数据字段计算"掩码 CRC32C"(masked crc32c)各占 4 字节; - 掩码算法在
_masked_crc32c(tfrecordio.py)中实现:(((crc >> 15) | (crc << 17)) + 0xa282ead8) & 0xffffffff,与 TensorFlow 的 TFRecord 格式保持兼容; - 读取时
read_record先读 12 字节头,校验长度掩码;再按长度读出数据与数据掩码,校验数据掩码;任一校验失败都会抛出ValueError(提示Not a valid TFRecord...);读到文件末尾返回None。
因此每条记录的实际开销是len(record) + 16字节(见encoded_num_bytes,tfrecordio.py),Beam 正是用它来推进读取偏移。
CRC32C 计算的底层实现依赖第三方库,tfrecordio.py 按python-snappy→google-crc32c→crcmod的顺序自动挑选最快的可用实现;如果三者都缺失,会抛出RuntimeError,提示执行pip install apache-beam[tfrecord](该 extras 会一并装齐依赖)。这意味着在不做任何配置的情况下,Beam 也能用纯 Python 的crcmod完成兼容性兜底,只是速度稍慢。
另外值得注意的是:_TFRecordSource在构造时传入splittable=False(tfrecordio.py),read_records还强制要求起始偏移为 0(tfrecordio.py),并从文件头逐条读取——这解释了为什么 TFRecord 源不可按偏移动态拆分,Beam 只能按文件粒度分发读取任务。
六、配套变换:WriteToTFRecord 与 ReadAllFromTFRecord
6.1 写入端:WriteToTFRecord
若要在管道里生成 TFRecord,使用WriteToTFRecord(tfrecordio.py):
from apache_beam.io.tfrecordio import WriteToTFRecord (p | beam.Create([b'foo', b'bar']) | WriteToTFRecord( file_path_prefix='/tmp/result', coder=coders.BytesCoder(), file_name_suffix='.tfrecord', num_shards=1, compression_type=CompressionTypes.GZIP))其关键参数:file_path_prefix为输出路径前缀(实际文件名为前缀 + 分片名 + 后缀);num_shards控制输出文件分片数;shard_name_template支持''、'-SSSSS-of-NNNNN'、'-W-SSSSS-of-NNNNN'、'-V-SSSSS-of-NNNNN'四种模板(W表示窗口区间,V表示 UTC 时间戳格式的窗口区间);流式(unbounded)管道会自动切换为窗口化的分片命名模板。测试TestWriteToTFRecord([tfrecordio_test.py](https://link.gitcode.com/i/c58a1048dfc8f2abadd70e1f51160fb1#L213-L241))验证了 GZIP 与 AUTO 压缩写入,并用tf.python_io.tf_record_iterator逐条读回比对。
6.2 动态文件列表:ReadAllFromTFRecord
当文件路径列表本身是运行时数据(例如先查询数据库或清单文件得到路径),用ReadAllFromTFRecord(tfrecordio.py):
from apache_beam.io.tfrecordio import ReadAllFromTFRecord (p | beam.Create(globs) # 每条元素是一个文件路径 | ReadAllFromTFRecord( coder=coders.BytesCoder(), compression_type=CompressionTypes.AUTO))它额外支持with_filename=True,此时输出从纯数据变为(文件名, 数据)的键值对,方便下游追踪每条记录来自哪个文件。对应测试见 tfrecordio_test.py:同一 glob 被Create成多条元素时,读取结果按元素数量成倍展开(如[b'foo', b'bar'] * 3)。
七、测试与工程实践佐证
仓库中的 tfrecordio_test.py 是理解本主题最好的补充材料:
TestTFRecordSink直接调用_TFRecordSink与_TFRecordUtil.write_record,逐字节断言写入结果(如b'foo'编码后的完整记录内容);TestReadFromTFRecord覆盖BytesCoder搭配AUTO/GZIP/DEFLATE读取;test_end2end(tfrecordio_test.py)演示了与本主题文档完全一致的实践路径:先用pickle.dump把随机矩阵序列化成字节、WriteToTFRecord落盘,再用ReadFromTFRecord(file_path_prefix + '-*')读回并断言相等——这正是"pickle 字节 ↔ TFRecord 文件"读写闭环的官方测试版。
从源码结构看,ReadFromTFRecord.expand最终执行pvalue.pipeline | Read(self._source)(tfrecordio.py),即把_TFRecordSource包装进 Beam 的Read变换,融入任何 Runner(Direct、Dataflow、Flink、Spark 等)的标准执行流程。
八、使用注意事项
- 依赖安装:读取/写入 TFRecord 依赖 CRC32C 计算库,建议
pip install apache-beam[tfrecord],或单独安装python-snappy/google-crc32c/crcmod之一,避免回退到慢速路径甚至报RuntimeError。 - pickle 安全:
pickle.loads反序列化不可信数据存在代码执行风险,仅应在数据来源可信(如自有训练数据、内部 ETL 产物)时使用;生产环境更推荐写入方使用tf.train.Example等显式 schema 的编码,读取端用协议解析器还原。 - 文件不可拆分:由于 TFRecord 源
splittable=False且必须从偏移 0 读取,超大文件应提前分片(多个文件),让 Beam 按文件并行读取;单文件内部无法并发切分。 - 路径校验时机:
validate=True在管道构建期检查文件存在性,适用于静态文件;文件由上游动态生成时应关闭该校验。
掌握ReadFromTFRecord的参数组合与底层 TFRecord 格式之后,你就能在 Beam 管道中无缝接入 TensorFlow 生态的数据文件,无论是离线训练样本抽取、跨 Runner 的数据迁移,还是结合WriteToTFRecord构建可复现的数据流水线。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Python SDK 实战:使用 ReadFromAvro 与 PipelineOptions 读取 Avro 文件
Apache Beam Python SDK 实战:使用 ReadFromAvro 与 PipelineOptions 读取 Avro 文件 导读 本文以 Ap
大数据批处理流处理数据工程page-agent路线图前瞻:send_keys、upload_file与表格解析三大功能何时到来?
page agent路线图前瞻:send_keys、upload_file与表格解析三大功能何时到来? Page Agent(page agent)是一个运行在
大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战:使用 TextIO 从文本文件读取 PCollection
Apache Beam Kotlin Kata 实战:使用 TextIO 从文本文件读取 PCollection 创建 Beam 管道时,最常见的第一步就是从文
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考