【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 通过一套统一编程模型同时支持批处理(Batch)与流处理(Streaming),而 I/O 连接器(I/O Connectors)正是打通 Beam 与外部数据系统的桥梁——它们以 Read/Write Transform 的形式,让管道可以用一致、可分布式扩展的方式读写各类存储系统。本文以 connectors.md 为骨架,完整梳理 Beam 内置连接器及其跨 SDK 支持矩阵、第三方连接器生态,并结合仓库源码(Java/Python/Go/TypeScript 各 SDK 与集成测试目录)讲清 native 实现与 X-language 跨语言实现的区别,帮助你在选型连接器、构建读写管道时快速定位到正确的用法与资源。
一、什么是 Apache Beam I/O 连接器
Apache Beam I/O 连接器为最主流的数据存储系统提供了读写 Transform,使 Beam 用户能够直接受益于原生优化的连接能力。借助这些 I/O,Beam 管道可以以统一且分布式的方式,从外部存储类型中读取数据、或向外部存储类型写入数据。
一个连接器通常由两部分组成:
- Source(数据源):把外部系统中的数据读入
PCollection; - Sink(数据汇):把
PCollection中的元素写入外部系统。
在 Beam 中,所有 Source 与 Sink 本质上都是复合 Transform(Composite Transform),它们可以被任意 Runner 解析并分布到多台工作节点上并行执行。连接器的能力可以用五个关键维度描述:是否支持 Source / Sink、在 Java / Python / Go / TypeScript 各 SDK 中的可用性、以及是否支持 Batch / Streaming 模式。这正是 connectors.md 中那张总览表的列结构。
native 实现与 X-language 实现
总览表中反复出现的native与via X-language两种标注,代表了连接器在某个语言 SDK 中接入的两种方式:
- native(原生实现):连接器的读写 Transform 由该 SDK 直接实现,用户无需额外启动跨语言服务即可使用。例如 Java 的
KafkaIO位于 sdks/java/io/kafka,Python 的ReadFromText位于 textio.py。 - via X-language(跨语言实现):该连接器实际由一个 SDK(通常是 Java)实现,其他 SDK 通过 Apache Beam 的 multi-language pipelines 框架 调用它。从源码可以清晰看到这一机制:例如 Python 的 Kafka 连接器 kafka.py 中,
ReadFromKafka是一个ExternalTransform,内部通过beam:transform:org.apache.beam:kafka_read_with_metadata:v1这类 URN 指向 Java 端实现;Go 侧则统一放在 sdks/go/pkg/beam/io/xlang 目录下(含 kafkaio、jdbcio、debeziumio、bigqueryio 等);TypeScript 侧则位于 sdks/typescript/src/apache_beam/io(如 kafka.ts、pubsub.ts、bigqueryio.ts)。
二、内置 I/O 连接器总览(完整矩阵)
下表是 connectors.md 的核心内容,按功能域分组成多张子表,保证每一行信息(Source/Sink 支持、各语言支持方式、Batch/Streaming 支持)完整保留。
阅读约定:✔ 表示支持;✘ 表示不支持;"Not available" 表示该 SDK 暂无此连接器;"via X-language" 表示通过跨语言管道框架调用。
2.1 文件与文件格式连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| FileIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | Not available | ✔ | ✔ |
| AvroIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✔ |
| TextIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✔ |
| TFRecordIO | ✔ | ✔ | ✔ native | ✔ native | Not available | Not available | ✔ | ✘ |
| XmlIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| TikaIO | ✔ | ✘ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| ParquetIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✘ |
| ThriftIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
这组连接器全部位于 Java SDK 的 sdks/java/io 目录(如 parquet、thrift、tika、xml、csv、json 等),Python 侧对应 sdks/python/apache_beam/io 下的 avroio.py、parquetio.py、tfrecordio.py、textio.py 等模块。其中 TextIO 和 ParquetIO 还配有专门的性能指标页与使用指南(见下文)。
2.2 文件系统连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| HadoopFileSystem | ✔ | ✔ | ✔ native | ✔ native | Not available | ✔ via X-language | ✔ | ✘ |
| GcsFileSystem | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✘ |
| LocalFileSystem | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✘ |
| S3FileSystem | ✔ | ✔ | ✔ native | ✔ native | Not available | ✔ via X-language | ✔ | ✘ |
| In-memory | ✘ | ✔ | ✘ | ✘ | ✔ native | ✘ | ✔ | ✘ |
文件系统连接器统一实现了 Beam 的FileSystem抽象,供 TextIO、AvroIO 等文件类连接器复用。从 Go 侧源码结构可以看到 filesystem 下分列 gcs、local、memfs(In-memory 的 Go 实现即 memfs)。注意:文件系统类连接器在表格中均标注 Streaming 为 ✘,即它们只服务于批处理场景的文件读写。
2.3 消息队列与事件流连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| KinesisIO | ✔ | ✔ | ✔ native | ✔ via X-language | Not available | Not available | ✔ | ✔ |
| AmqpIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| KafkaIO | ✔ | ✔ | ✔ native | ✔ via X-language | ✔ via X-language | ✔ via X-language | ✔ | ✔ |
| PubSubIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | ✔ via X-language | ✔ | ✔ |
| JmsIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| MqttIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| RabbitMqIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| SqsIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| SnsIO | ✘ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
这组连接器是流式管道的核心。其中 KafkaIO 与 PubSubIO 是四语言覆盖最完整的代表:KafkaIO 在 Java 原生实现,Python/Go/TypeScript 均通过 X-language 接入(对应 kafka.py、xlang/kafkaio 与 kafka.ts);PubSubIO 则在 Java、Python、Go 三个 SDK 都有原生实现,TypeScript 走 X-language。
2.4 关系型与数据库连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| JdbcIO | ✔ | ✔ | ✔ native | ✔ via X-language | ✔ via X-language | Not available | ✔ | ✘ |
| CassandraIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| HBaseIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| KuduIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| DatabaseIO | ✔ | ✔ | ✘ | ✘ | ✔ native | Not available | ✔ | ✘ |
| MongoDbIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | Not available | ✔ | ✘ |
| MongoDbGridFSIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| ClickHouseIO | ✘ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| RedisIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| Neo4j | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| SingleStoreDB | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
JdbcIO 是数据库连接器的通用入口(MySQL、PostgreSQL 等),其实现位于 sdks/java/io/jdbc。Python 的apache_beam.io.jdbc与 Go 的 xlang/jdbcio 均通过跨语言方式复用 Java 实现。
2.5 大数据生态(Hadoop / 搜索 / 时序)连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| HadoopFormatIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| HCatalogIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| SolrIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| ElasticsearchIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| InfluxDB | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| SplunkIO | ✘ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
2.6 Google Cloud 服务连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| BigQueryIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native + via X-language | ✔ via X-language | ✔ | ✔ |
| BigTableIO | ✔ | ✔ | ✔ native | ✔ native(sink) + via X-language | ✔ native(sink) + via X-language | Not available | ✔ | ✔ |
| DatastoreIO | ✔ | ✔ | ✔ native | ✔ native | ✔ native | Not available | ✔ | ✔ |
| SpannerIO | ✔ | ✔ | ✔ native | ✔ via X-language | ✔ native | Not available | ✔ | ✔ |
| Pub/Sub Lite | ✔ | ✔ | ✔ native | ✔ via X-language | Not available | ✔ via X-language | ✔ | ✔ |
| Firestore IO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| FhirIO | ✔ | ✔ | ✔ native | Not available | ✔ native | Not available | ✔ | ✔ |
| HL7v2IO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| DicomIO | ✔ | ✔ | ✔ native | ✔ native | Not available | Not available | ✔ | ✔ |
| GoogleAdsIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
Google Cloud 连接器统一实现在 sdks/java/io/google-cloud-platform(含 bigquery、bigtable、datastore、spanner、pubsub、healthcare(Fhir/HL7v2/Dicom)、firestore、pubsublite 等子包),Python 侧集中在 sdks/python/apache_beam/io/gcp。
2.7 云平台与消息总线(Amazon / 其他)连接器
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| DynamoDBIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| DebeziumIO | ✔ | ✘ | ✔ native | ✔ via X-language | ✔ via X-language | Not available | ✔ | ✔ |
| SnowflakeIO | ✔ | ✔ | ✔ native | ✔ via X-language | Not available | Not available | ✔ | ✘ |
| CdapIO | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| FlinkStreaming ImpulseSource | ✔ | ✘ | Not available | ✔ native | Not available | Not available | ✔ | ✔ |
| SparkReceiverIO | ✔ | ✘ | ✔ native | Not available | Not available | Not available | ✘ | ✔ |
| GenerateSequence | ✔ | ✘ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
其中GenerateSequence是一个特殊的测试/仿真连接器,用于生成递增的整数序列,可同时用于批处理与流式场景(如驱动窗口计算演示)。SparkReceiverIO是唯一在 Batch 列为 ✘ 而 Streaming 为 ✔ 的连接器,因为它专门用于接收 Spark Streaming 的数据。
三、如何解读连接器矩阵:列与标注的实战含义
Source / Sink 两列:决定连接器能"读"还是能"写"。例如 SnsIO、SplunkIO、ClickHouseIO、Tinybird 等只支持 Sink(写入);TikaIO、DebeziumIO、SparkReceiverIO、GenerateSequence 等只支持 Source(读取)。设计管道前先确认方向,避免引入方向不支持的连接器。
语言列:
native表示该 SDK 直接实现,性能与调试体验最佳,且通常文档示例最全;via X-language表示依赖跨语言管道,使用前需要配置相应的 Expansion Service / 端口转发,部署复杂度略高。对于同一连接器,优先选择 native 实现的语言。Batch / Streaming 两列:决定连接器适用的管道模式。例如 TFRecordIO、ParquetIO、XmlIO、ThriftIO 及各文件系统连接器都只支持 Batch;而 KafkaIO、PubSubIO、KinesisIO、JmsIO、MqttIO、RabbitMqIO、SqsIO、AmqpIO、SolrIO、SplunkIO、InfluxDB、DebeziumIO、HadoopFormatIO 等同时支持 Batch 与 Streaming,可灵活用于 Lambda 架构或流批一体管道。
性能指标与专属指南:矩阵中部分连接器名称旁带
metrics或guide标记,例如 TextIO 与 BigQueryIO 的 metrics 链接指向性能指标页(仓库中见 website/www/site/content/en/performance/textio/_index.md 等),ParquetIO、HadoopFormatIO、HCatalogIO、SnowflakeIO、CdapIO、SparkReceiverIO、SingleStoreDB 等则配有独立的使用指南页(位于 website/www/site/content/en/documentation/io/built-in 目录)。配置具体参数前,建议先阅读对应指南。
四、关键连接器的源码级解读
4.1 KafkaIO:流式消息连接器的范式
KafkaIO.java 是 Java SDK 中最具代表性的流式连接器,其 API 结构体现了 Beam 连接器的通用设计约定(该约定本身被文档化为 io-standards.md):
- 顶层静态方法:
KafkaIO.readBytes()(读取裸字节)、KafkaIO.read()(读取泛型 K/V)、KafkaIO.write(),用户从这三个入口开始构建读写; - 内部
Read<K, V>与Write<K, V>抽象类采用流式 Builder 风格配置:如getTopics()、getTopicPartitions()、getConsumerConfig()、getWatermarkFn()、getTimestampPolicyFactory()等,分别对应主题、分区、消费者参数、水位线函数与时间戳策略; - 支持
ReadFromKafkaViaSDF与ReadFromKafkaViaUnbounded两条读取路径——前者基于 Splittable DoFn,后者基于传统 UnboundedSource,二者都封装在统一的Read接口之下。
从 developing-io-overview.md 可以确认,Splittable DoFn 是 Beam 官方推荐的 Source 实现方式(支持 checkpoint、水位线控制、backlog 追踪、动态分片),Kafka、Parquet、HL7v2 等连接器均为其典型实践。
4.2 TextIO:文件读写的最简实践
TextIO.java 位于 Java 核心 SDK(sdks/java/core),提供TextIO.read()/TextIO.write()以及ReadAll、ReadFiles变体;Python 侧对应 textio.py 中的ReadFromText、ReadFromTextWithFilename与WriteToText;Go 侧对应 sdks/go/pkg/beam/io/textio 包。由于 TextIO 依赖FileSystems抽象,它可以无缝读写本地文件、GCS、S3、HDFS 等多种文件系统,是理解 Beam 文件类连接器的最佳入门样例。
4.3 连接器在仓库中的代码落位
- Java:所有内置连接器源码位于 sdks/java/io,按系统分目录(kafka、kinesis、jdbc、mongodb、parquet、cassandra、hbase、hcatalog、solr、elasticsearch、google-cloud-platform、amazon-web-services2 等);
- Python:位于 sdks/python/apache_beam/io,文件命名与类一一对应(如 avroio.py、parquetio.py、mongodbio.py、kinesis.py、snowflake.py、debezium.py);
- Go:位于 sdks/go/pkg/beam/io,原生实现在顶层各包(textio、avroio、parquetio、bigqueryio、bigtableio、pubsubio、mongodbio、spannerio 等),跨语言实现在 xlang 子目录(kafkaio、jdbcio、debeziumio、bigqueryio、bigtableio、schemaio);
- TypeScript:位于 sdks/typescript/src/apache_beam/io,目前提供 avroio、bigqueryio、kafka、parquetio、pubsub、pubsublite、textio、schemaio 等,全部走 X-language。
五、内置连接器的质量保障:集成测试与性能基准
连接器进入"内置"列表并非仅靠代码提交,还配套了严格的测试与性能跟踪体系:
- 集成测试(Integration Tests):仓库的 it 目录为多数内置连接器提供了面向真实存储实例的集成测试,覆盖 cassandra、elasticsearch、google-cloud-platform、jdbc、kafka、mongodb、neo4j、splunk 等。例如 Kafka 的资源管理辅助类 KafkaResourceManager.java 负责在集成测试中创建/销毁 Kafka 集群资源。测试方法论见 testing.md:单元测试用内存版/伪造的数据存储验证核心行为,集成测试用真实实例(千行到数十 GB 规模)捕捉多 Worker 并发读写、数据副本等边界问题,且不做单独的基准测试,而是通过可参数化的集成测试覆盖性能场景。
- 性能指标(Performance Metrics):Google Cloud Dataflow 团队使用 Dataflow Runner 定期运行各内置连接器的性能测试并公开指标(如 TextIO 性能页),这正是总览表中 TextIO、BigQueryIO、BigTableIO 名称旁 metrics 标记的来源。
- 连接器规范(I/O Standards):io-standards.md 定义了内置连接器的开发规范:Java 主类统一命名为
{Connector}IO、置于org.apache.beam.sdk.io.{connector}包;Python 统一为apache_beam.io.{connector}包并定义__all__;文档需包含 Before you start、Supported Features(关系型特性表)、Authentication、Limitations 等固定小节,并更新本总览表。
六、其他(非内置)I/O 连接器
除 Beam 仓库内的内置连接器外,社区还维护了一批托管在各自项目仓库中的第三方连接器,同样以 Source/Sink 形式接入 Beam。下表完整继承自 connectors.md:
| 连接器 | Source | Sink | Java | Python | Go | TypeScript | Batch | Streaming |
|---|---|---|---|---|---|---|---|---|
| Solace | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| SAP Hana to Google BigQuery | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| MySQL | ✔ | ✔ | Not available | ✔ native | Not available | Not available | ✔ | ✘ |
| TrepWsIO | ✔ | ✘ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| KineticaDB | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
| Cognite Data Fusion | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✔ |
| Pyodbc | ✔ | ✔ | Not available | ✔ native | Not available | Not available | ✔ | ✘ |
| Go Connect | ✔ | ✔ | ✘ | ✘ | ✔ native | Not available | ✔ | ✔ |
| Tinybird | ✘ | ✔ | Not available | ✔ native | Not available | Not available | ✔ | ✔ |
| Cloud SQL | ✔ | ✘ | Not available | ✔ native | Not available | Not available | ✔ | ✘ |
| Cloud Bigtable (HBase based) | ✔ | ✔ | ✔ native | Not available | Not available | Not available | ✔ | ✘ |
这类连接器的接入方式(如 Maven 坐标、Python 包安装、配置参数)以其各自项目仓库为准。如果你的团队需要接入列表之外的数据存储,可以按照 developing-io-overview.md 与 developing-io-java.md 自行开发自定义连接器,并遵循 io-standards.md 的规范与 testing.md 的测试要求,待成熟后按贡献流程将其纳入内置连接器列表。
七、连接器选型速查建议
结合以上矩阵,给出四条可操作的选型准则:
- 先看语言列选 native:在 Java 中优先用 KafkaIO、BigQueryIO、SpannerIO 等原生实现;在 Python 中优先用 ReadFromText、ReadFromBigQuery、ReadFromPubSub 等原生 Transform;只有目标连接器在所选语言中标注 via X-language 时才配置跨语言服务。
- 再按模式过滤 Batch/Streaming:纯批管道避免引入仅 Streaming 的连接器(如 SparkReceiverIO);纯流管道注意 TFRecordIO、ParquetIO、JdbcIO 等仅支持 Batch 的选项,这类连接器只能用于流式管道的批式落地步骤。
- 关注带指南/指标标记的连接器:使用 ParquetIO、HadoopFormatIO、HCatalogIO、SnowflakeIO、CdapIO、SparkReceiverIO、SingleStoreDB 前先读对应指南页(built-in 目录);对 TextIO、BigQueryIO、BigTableIO 可查阅性能指标页评估规模。
- 同一连接器多语言能力可互补:例如 BigQueryIO 在 Go 中同时提供 native 与 X-language 两套入口,可根据是否已部署 Java Expansion Service 灵活选择;BigTableIO 在 Python/Go 中 native 仅支持 Sink(写入),需要读 Bigtable 时建议使用 Java 原生实现或 X-language 通道。
Apache Beam 的连接器生态覆盖文件、消息、数据库、云服务与大数据组件,本文的矩阵与源码落位信息可以直接作为你在 connectors.md 之外快速检索实现、定位测试与撰写管道代码的起点。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam I/O Connectors 全览:内置连接器矩阵、跨语言机制与 Managed API 实践
Apache Beam I/O Connectors 全览:内置连接器矩阵、跨语言机制与 Managed API 实践 Apache Beam 的 I/O 连接
大数据批处理流处理数据工程Apache Beam Snowflake I/O 连接器深度解析:认证配置、COPY 读写原理与 Python 跨语言支持
Apache Beam Snowflake I/O 连接器深度解析:认证配置、COPY 读写原理与 Python 跨语言支持 本文基于 Apache Beam
大数据批处理流处理数据工程Apache Beam Managed I/O 连接器完全指南:统一配置接口与 Runner 托管实践
Apache Beam Managed I/O 连接器完全指南:统一配置接口与 Runner 托管实践 Apache Beam 的 Managed API 以一
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考