☰
Apache Beam I/O 连接器全览:内置连接器矩阵、跨语言支持与选型实践
2026/10/12 5:59:37 网站建设 项目流程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

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 文件与文件格式连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
FileIO✔✔✔ native✔ native✔ nativeNot available✔✔
AvroIO✔✔✔ native✔ native✔ native✔ via X-language✔✔
TextIO✔✔✔ native✔ native✔ native✔ via X-language✔✔
TFRecordIO✔✔✔ native✔ nativeNot availableNot available✔✘
XmlIO✔✔✔ nativeNot availableNot availableNot available✔✘
TikaIO✔✘✔ nativeNot availableNot availableNot available✔✔
ParquetIO✔✔✔ native✔ native✔ native✔ via X-language✔✘
ThriftIO✔✔✔ nativeNot availableNot availableNot 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 文件系统连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
HadoopFileSystem✔✔✔ native✔ nativeNot available✔ via X-language✔✘
GcsFileSystem✔✔✔ native✔ native✔ native✔ via X-language✔✘
LocalFileSystem✔✔✔ native✔ native✔ native✔ via X-language✔✘
S3FileSystem✔✔✔ native✔ nativeNot available✔ via X-language✔✘
In-memory✘✔✘✘✔ native✘✔✘

文件系统连接器统一实现了 Beam 的FileSystem抽象,供 TextIO、AvroIO 等文件类连接器复用。从 Go 侧源码结构可以看到 filesystem 下分列 gcs、local、memfs(In-memory 的 Go 实现即 memfs)。注意:文件系统类连接器在表格中均标注 Streaming 为 ✘,即它们只服务于批处理场景的文件读写。

2.3 消息队列与事件流连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
KinesisIO✔✔✔ native✔ via X-languageNot availableNot available✔✔
AmqpIO✔✔✔ nativeNot availableNot availableNot available✔✔
KafkaIO✔✔✔ native✔ via X-language✔ via X-language✔ via X-language✔✔
PubSubIO✔✔✔ native✔ native✔ native✔ via X-language✔✔
JmsIO✔✔✔ nativeNot availableNot availableNot available✔✔
MqttIO✔✔✔ nativeNot availableNot availableNot available✔✔
RabbitMqIO✔✔✔ nativeNot availableNot availableNot available✔✔
SqsIO✔✔✔ nativeNot availableNot availableNot available✔✔
SnsIO✘✔✔ nativeNot availableNot availableNot 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 关系型与数据库连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
JdbcIO✔✔✔ native✔ via X-language✔ via X-languageNot available✔✘
CassandraIO✔✔✔ nativeNot availableNot availableNot available✔✘
HBaseIO✔✔✔ nativeNot availableNot availableNot available✔✘
KuduIO✔✔✔ nativeNot availableNot availableNot available✔✘
DatabaseIO✔✔✘✘✔ nativeNot available✔✘
MongoDbIO✔✔✔ native✔ native✔ nativeNot available✔✘
MongoDbGridFSIO✔✔✔ nativeNot availableNot availableNot available✔✘
ClickHouseIO✘✔✔ nativeNot availableNot availableNot available✔✘
RedisIO✔✔✔ nativeNot availableNot availableNot available✔✘
Neo4j✔✔✔ nativeNot availableNot availableNot available✔✘
SingleStoreDB✔✔✔ nativeNot availableNot availableNot available✔✘

JdbcIO 是数据库连接器的通用入口(MySQL、PostgreSQL 等),其实现位于 sdks/java/io/jdbc。Python 的apache_beam.io.jdbc与 Go 的 xlang/jdbcio 均通过跨语言方式复用 Java 实现。

2.5 大数据生态(Hadoop / 搜索 / 时序)连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
HadoopFormatIO✔✔✔ nativeNot availableNot availableNot available✔✔
HCatalogIO✔✔✔ nativeNot availableNot availableNot available✔✔
SolrIO✔✔✔ nativeNot availableNot availableNot available✔✔
ElasticsearchIO✔✔✔ nativeNot availableNot availableNot available✔✔
InfluxDB✔✔✔ nativeNot availableNot availableNot available✔✔
SplunkIO✘✔✔ nativeNot availableNot availableNot available✔✔

2.6 Google Cloud 服务连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
BigQueryIO✔✔✔ native✔ native✔ native + via X-language✔ via X-language✔✔
BigTableIO✔✔✔ native✔ native(sink) + via X-language✔ native(sink) + via X-languageNot available✔✔
DatastoreIO✔✔✔ native✔ native✔ nativeNot available✔✔
SpannerIO✔✔✔ native✔ via X-language✔ nativeNot available✔✔
Pub/Sub Lite✔✔✔ native✔ via X-languageNot available✔ via X-language✔✔
Firestore IO✔✔✔ nativeNot availableNot availableNot available✔✘
FhirIO✔✔✔ nativeNot available✔ nativeNot available✔✔
HL7v2IO✔✔✔ nativeNot availableNot availableNot available✔✔
DicomIO✔✔✔ native✔ nativeNot availableNot available✔✔
GoogleAdsIO✔✔✔ nativeNot availableNot availableNot 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 / 其他)连接器

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
DynamoDBIO✔✔✔ nativeNot availableNot availableNot available✔✔
DebeziumIO✔✘✔ native✔ via X-language✔ via X-languageNot available✔✔
SnowflakeIO✔✔✔ native✔ via X-languageNot availableNot available✔✘
CdapIO✔✔✔ nativeNot availableNot availableNot available✔✔
FlinkStreaming ImpulseSource✔✘Not available✔ nativeNot availableNot available✔✔
SparkReceiverIO✔✘✔ nativeNot availableNot availableNot available✘✔
GenerateSequence✔✘✔ nativeNot availableNot availableNot available✔✔

其中GenerateSequence是一个特殊的测试/仿真连接器,用于生成递增的整数序列,可同时用于批处理与流式场景(如驱动窗口计算演示)。SparkReceiverIO是唯一在 Batch 列为 ✘ 而 Streaming 为 ✔ 的连接器,因为它专门用于接收 Spark Streaming 的数据。

三、如何解读连接器矩阵:列与标注的实战含义

  1. Source / Sink 两列:决定连接器能"读"还是能"写"。例如 SnsIO、SplunkIO、ClickHouseIO、Tinybird 等只支持 Sink(写入);TikaIO、DebeziumIO、SparkReceiverIO、GenerateSequence 等只支持 Source(读取)。设计管道前先确认方向,避免引入方向不支持的连接器。

  2. 语言列:native表示该 SDK 直接实现,性能与调试体验最佳,且通常文档示例最全;via X-language表示依赖跨语言管道,使用前需要配置相应的 Expansion Service / 端口转发,部署复杂度略高。对于同一连接器,优先选择 native 实现的语言。

  3. Batch / Streaming 两列:决定连接器适用的管道模式。例如 TFRecordIO、ParquetIO、XmlIO、ThriftIO 及各文件系统连接器都只支持 Batch;而 KafkaIO、PubSubIO、KinesisIO、JmsIO、MqttIO、RabbitMqIO、SqsIO、AmqpIO、SolrIO、SplunkIO、InfluxDB、DebeziumIO、HadoopFormatIO 等同时支持 Batch 与 Streaming,可灵活用于 Lambda 架构或流批一体管道。

  4. 性能指标与专属指南:矩阵中部分连接器名称旁带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:

连接器SourceSinkJavaPythonGoTypeScriptBatchStreaming
Solace✔✔✔ nativeNot availableNot availableNot available✔✔
SAP Hana to Google BigQuery✔✔✔ nativeNot availableNot availableNot available✔✘
MySQL✔✔Not available✔ nativeNot availableNot available✔✘
TrepWsIO✔✘✔ nativeNot availableNot availableNot available✔✔
KineticaDB✔✔✔ nativeNot availableNot availableNot available✔✘
Cognite Data Fusion✔✔✔ nativeNot availableNot availableNot available✔✔
Pyodbc✔✔Not available✔ nativeNot availableNot available✔✘
Go Connect✔✔✘✘✔ nativeNot available✔✔
Tinybird✘✔Not available✔ nativeNot availableNot available✔✔
Cloud SQL✔✘Not available✔ nativeNot availableNot available✔✘
Cloud Bigtable (HBase based)✔✔✔ nativeNot availableNot availableNot available✔✘

这类连接器的接入方式(如 Maven 坐标、Python 包安装、配置参数)以其各自项目仓库为准。如果你的团队需要接入列表之外的数据存储,可以按照 developing-io-overview.md 与 developing-io-java.md 自行开发自定义连接器,并遵循 io-standards.md 的规范与 testing.md 的测试要求,待成熟后按贡献流程将其纳入内置连接器列表。

七、连接器选型速查建议

结合以上矩阵,给出四条可操作的选型准则:

  1. 先看语言列选 native:在 Java 中优先用 KafkaIO、BigQueryIO、SpannerIO 等原生实现;在 Python 中优先用 ReadFromText、ReadFromBigQuery、ReadFromPubSub 等原生 Transform;只有目标连接器在所选语言中标注 via X-language 时才配置跨语言服务。
  2. 再按模式过滤 Batch/Streaming:纯批管道避免引入仅 Streaming 的连接器(如 SparkReceiverIO);纯流管道注意 TFRecordIO、ParquetIO、JdbcIO 等仅支持 Batch 的选项,这类连接器只能用于流式管道的批式落地步骤。
  3. 关注带指南/指标标记的连接器:使用 ParquetIO、HadoopFormatIO、HCatalogIO、SnowflakeIO、CdapIO、SparkReceiverIO、SingleStoreDB 前先读对应指南页(built-in 目录);对 TextIO、BigQueryIO、BigTableIO 可查阅性能指标页评估规模。
  4. 同一连接器多语言能力可互补:例如 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.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询