使用 Spark
Hudi Streamer
Hudi Streamer(包含在 hudi-utilities-slim-bundle 和 hudi-utilities-bundle 中)提供了从 DFS 或 Kafka 等不同数据源摄入数据的方式,并具备以下能力。
- 从 Kafka 精确一次(exactly once)地摄入新事件,从 Sqoop 执行增量导入,摄入
HiveIncrementalPuller的输出结果,或摄入 DFS 目录下的文件 - 支持传入数据使用 json、avro 或自定义的记录类型
- 管理检查点、回滚与恢复
- 从 DFS 或 Confluent schema registry 获取 Avro schema
- 支持插入转换(transformation)逻辑
重要
以下类已被重命名并迁移到 org.apache.hudi.utilities.streamer 包中。
DeltastreamerMultiWriterCkptUpdateFunc重命名为StreamerMultiWriterCkptUpdateFuncDeltaSync重命名为StreamSyncHoodieDeltaStreamer重命名为HoodieStreamerHoodieDeltaStreamerMetrics重命名为HoodieStreamerMetricsHoodieMultiTableDeltaStreamer重命名为HoodieMultiTableStreamer
为保持向后兼容,原始类仍然保留在 org.apache.hudi.utilities.deltastreamer 包中,但已被标记为弃用(deprecated)。
选项
展开此处可查看 Hudi Streamer 的 --help 输出,其中更详细地描述了它的各项能力。
[hoodie]$ spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer `ls packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle-*.jar` --help
Usage: <main class> [options]
Options:
--allow-commit-on-no-checkpoint-change
allow commits even if checkpoint has not changed before and after fetch
datafrom source. This might be useful in sources like SqlSource where
there is not checkpoint. And is not recommended to enable in continuous
mode.
Default: false
--base-file-format
File format for the base files. PARQUET (or) HFILE
Default: PARQUET
--bootstrap-index-class
subclass of BootstrapIndex
Default: org.apache.hudi.common.bootstrap.index.HFileBootstrapIndex
--bootstrap-overwrite
Overwrite existing target table, default false
Default: false
--checkpoint
Resume Hudi Streamer from this checkpoint.
--cluster-scheduling-minshare
Minshare for clustering as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 0
--cluster-scheduling-weight
Scheduling weight for clustering as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 1
--commit-on-errors
Commit even when some records failed to be written
Default: false
--compact-scheduling-minshare
Minshare for compaction as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 0
--compact-scheduling-weight
Scheduling weight for compaction as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 1
--config-hot-update-strategy-class
Configuration hot update in continuous mode
Default: <empty string>
--continuous
Hudi Streamer runs in continuous mode running source-fetch -> Transform
-> Hudi Write in loop
Default: false
--delta-sync-scheduling-minshare
Minshare for delta sync as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 0
--delta-sync-scheduling-weight
Scheduling weight for delta sync as defined in
https://spark.apache.org/docs/latest/job-scheduling.html
Default: 1
--disable-compaction
Compaction is enabled for MoR table by default. This flag disables it
Default: false
--enable-hive-sync
Enable syncing to hive
Default: false
--enable-sync
Enable syncing meta
Default: false
--filter-dupes
Should duplicate records from source be dropped/filtered out before
insert/bulk-insert
Default: false
--force-empty-sync
Force syncing meta even on empty commit
Default: false
--help, -h
--hoodie-conf
Any configuration that can be set in the properties file (using the CLI
parameter "--props") can also be passed command line using this
parameter. This can be repeated
Default: []
--ingestion-metrics-class
Ingestion metrics class for reporting metrics during ingestion
lifecycles.
Default: org.apache.hudi.utilities.streamer.HoodieStreamerMetrics
--initial-checkpoint-provider
subclass of
org.apache.hudi.utilities.checkpointing.InitialCheckpointProvider.
Generate check point for Hudi Streamer for the first run. This field
will override the checkpoint of last commit using the checkpoint field.
Use this field only when switching source, for example, from DFS source
to Kafka Source.
--max-pending-clustering
Maximum number of outstanding inflight/requested clustering. Delta Sync
will not happen unlessoutstanding clustering is less than this number
Default: 5
--max-pending-compactions
Maximum number of outstanding inflight/requested compactions. Delta Sync
will not happen unlessoutstanding compactions is less than this number
Default: 5
--max-retry-count
the max retry count if --retry-on-source-failures is enabled
Default: 3
--min-sync-interval-seconds
the min sync interval of each sync in continuous mode
Default: 0
--op
Takes one of these values : UPSERT (default), INSERT, BULK_INSERT,
INSERT_OVERWRITE, INSERT_OVERWRITE_TABLE, DELETE_PARTITION, DELETE
(DELETE extracts HoodieKeys from source records and deletes the
corresponding records from the table.)
Default: UPSERT
Possible Values: [INSERT, INSERT_PREPPED, UPSERT, UPSERT_PREPPED, BULK_INSERT, BULK_INSERT_PREPPED, DELETE, DELETE_PREPPED, BOOTSTRAP, INSERT_OVERWRITE, CLUSTER, DELETE_PARTITION, INSERT_OVERWRITE_TABLE, COMPACT, INDEX, ALTER_SCHEMA, LOG_COMPACT, UNKNOWN]
--payload-class
subclass of HoodieRecordPayload, that works off a GenericRecord.
Implement your own, if you want to do something other than overwriting
existing value
Default: org.apache.hudi.common.model.OverwriteWithLatestAvroPayload
--post-write-termination-strategy-class
Post writer termination strategy class to gracefully shutdown
Hudi Streamer in continuous mode
Default: <empty string>
--props
path to properties file on localfs or dfs, with configurations for
hoodie client, schema provider, key generator and data source. For
hoodie client props, sane defaults are used, but recommend use to
provide basic things like metrics endpoints, hive configs etc. For
sources, referto individual classes, for supported properties.
Properties in this file can be overridden by "--hoodie-conf"
Default: file:///Users/shiyanxu/src/test/resources/streamer-config/dfs-source.properties
--retry-interval-seconds
the retry interval for source failures if --retry-on-source-failures is
enabled
Default: 30
--retry-last-pending-inline-clustering, -rc
Retry last pending inline clustering plan before writing to sink.
Default: false
--retry-last-pending-inline-compaction
Retry last pending inline compaction plan before writing to sink.
Default: false
--retry-on-source-failures
Retry on any source failures
Default: false
--run-bootstrap
Run bootstrap if bootstrap index is not found
Default: false
--schemaprovider-class
subclass of org.apache.hudi.utilities.schema.SchemaProvider to attach
schemas to input & target table data, built in options:
org.apache.hudi.utilities.schema.FilebasedSchemaProvider.Source (See
org.apache.hudi.utilities.sources.Source) implementation can implement
their own SchemaProvider. For Sources that return Dataset<Row>, the
schema is obtained implicitly. However, this CLI option allows
overriding the schemaprovider returned by Source.
--source-class
Subclass of org.apache.hudi.utilities.sources to read data. Built-in
options: org.apache.hudi.utilities.sources.{JsonDFSSource (default),
AvroDFSSource, JsonKafkaSource, AvroKafkaSource, HiveIncrPullSource}
Default: org.apache.hudi.utilities.sources.JsonDFSSource
--source-limit
Maximum amount of data to read from source. Default: No limit, e.g:
DFS-Source => max bytes to read, Kafka-Source => max events to read
Default: 9223372036854775807
--source-ordering-field
Field within source record to decide how to break ties between records
with same key in input data. Default: 'ts' holding unix timestamp of
record
Default: ts
--spark-master
spark master to use, if not defined inherits from your environment
taking into account Spark Configuration priority rules (e.g. not using
spark-submit command).
Default: <empty string>
--sync-tool-classes
Meta sync client tool, using comma to separate multi tools
Default: org.apache.hudi.hive.HiveSyncTool
* --table-type
Type of table. COPY_ON_WRITE (or) MERGE_ON_READ
* --target-base-path
base path for the target hoodie table. (Will be created if did not exist
first time around. If exists, expected to be a hoodie table)
* --target-table
name of the target table
--transformer-class
A subclass or a list of subclasses of
org.apache.hudi.utilities.transform.Transformer. Allows transforming raw
source Dataset to a target Dataset (conforming to target schema) before
writing. Default : Not set. E.g. -
org.apache.hudi.utilities.transform.SqlQueryBasedTransformer (which
allows a SQL query templated to be passed as a transformation function).
Pass a comma-separated list of subclass names to chain the
transformations. If there are two or more transformers using the same
config keys and expect different values for those keys, then transformer
can include an identifier. E.g. -
tr1:org.apache.hudi.utilities.transform.SqlQueryBasedTransformer. Here
the identifier tr1 can be used along with property key like
`hoodie.streamer.transformer.sql.tr1` to identify properties related to
the transformer. So effective value for
`hoodie.streamer.transformer.sql` is determined by key
`hoodie.streamer.transformer.sql.tr1` for this transformer. If
identifier is used, it should be specified for all the transformers.
Further the order in which transformer is applied is determined by the
occurrence of transformer irrespective of the identifier used for the
transformer. For example: In the configured value below tr2:org.apache.hudi.utilities.transform.SqlQueryBasedTransformer,tr1:org.apache.hudi.utilities.transform.SqlQueryBasedTransformer
, tr2 is applied before tr1 based on order of occurrence.该工具接收一个分层组合的属性文件,并提供可插拔的接口用于数据提取、键生成和提供 Schema。从 Kafka 和 DFS 摄取数据的示例配置位于 hudi-utilities/src/test/resources/streamer-config 目录下。
例如:一旦你启动并运行了 Confluent Kafka 和 Schema Registry,就可以使用 schema-registry 仓库提供的 (impressions.avro) 生成一些测试数据
[confluent-5.0.0]$ bin/ksql-datagen schema=../impressions.avro format=avro topic=impressions key=impressionid然后按如下方式将其摄取。
[hoodie]$ spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer `ls packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle-*.jar` \
--props file://${PWD}/hudi-utilities/src/test/resources/streamer-config/kafka-source.properties \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.AvroKafkaSource \
--source-ordering-field impresssiontime \
--target-base-path file:\/\/\/tmp/hudi-streamer-op \
--target-table uber.impressions \
--op BULK_INSERT在某些情况下,你可能希望提前将现有表迁移到 Hudi 中。请参阅迁移指南。
使用 hudi-utilities-slim-bundle bundle jar
推荐使用 hudi-utilities-slim-bundle,并将其与对应所使用 Spark 版本的 Hudi Spark bundle 一起使用,才能让工具基于 Spark 运行,例如:--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0。
并发控制
通过 Hudi Streamer 使用乐观并发控制(OCC),需要在传递给作业的属性文件中配置以下参数。
hoodie.write.concurrency.mode=optimistic_concurrency_control
hoodie.write.lock.provider=<lock-provider-classname>
hoodie.cleaner.policy.failed.writes=LAZY例如,将这些配置添加到 kafka-source.properties 文件中并传递给 Hudi Streamer,即可启用 OCC。随后可以按如下方式触发 Hudi Streamer 作业:
[hoodie]$ spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer `ls packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle-*.jar` \
--props file://${PWD}/hudi-utilities/src/test/resources/streamer-config/kafka-source.properties \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.AvroKafkaSource \
--source-ordering-field impresssiontime \
--target-base-path file:///tmp/hudi-streamer-op \
--target-table uber.impressions \
--op BULK_INSERT更多关于并发控制的内容,请参见并发控制概念章节。
检查点(Checkpointing)
Hudi Streamer 使用检查点来记录已经读取过的数据,从而无需重新处理全部数据即可恢复运行。使用 Kafka 源时,检查点是 Kafka Offset;使用 DFS 源时,检查点是最新读取文件的"最后修改"时间戳。检查点保存在 .hoodie 提交文件中。完成时间(completion time)检查点存储在 streamer.checkpoint.key.v2 下,请求时间(request time)检查点存储在 deltastreamer.checkpoint.key 下。具体使用哪一个取决于数据源,对于 Hudi 增量源,还取决于表版本——详见下文。
如果需要更改检查点以重新处理或重放数据,可以使用以下选项:
--checkpoint会将提交文件中对应的重置键设置为当前检查点的覆盖值,即streamer.checkpoint.reset.key.v2或deltastreamer.checkpoint.reset_key。检查点的格式取决于 KAFKA_CHECKPOINT_TYPE。默认情况下(类型为string),检查点应以如下格式提供:topicName,0:offset0,1:offset1,2:offset2。类型为timestamp时,检查点应提供为目标时间戳的 long 值。类型为single_offset时,我们假设主题只包含一个分区,因此检查点应提供为目标偏移量的 long 值。--source-limit会设置从源读取的最大数据量。对于 DFS 源,这是读取的最大字节数;对于 Kafka,这是读取的最大事件数。
重置 Hudi 增量源的检查点
当 Hudi Streamer 使用 HoodieIncrSource 将目标表写入表版本 8 或更高版本时,它通过完成时间(而非请求的即时时间)来跟踪进度,并将其记录在提交元数据的 streamer.checkpoint.key.v2 下。裸时间戳在这两者之间存在歧义,因此 --checkpoint 必须指明它属于哪一种。这种格式从 Hudi 1.0.1 起为必需;在 1.0.0 上,前缀不会被识别,--checkpoint 接受一个裸的完成时间。
# resume after the instant with this requested instant time
--checkpoint resumeFromInstantRequestTime:20250110120000000
# resume after the instant with this completion time
--checkpoint resumeFromInstantCompletionTime:20250110120005000两种形式都会从同一个位置恢复。请求时间形式会在内部被转换为该时刻的完成时间,无论采用哪种方式,摄入都会按完成时间顺序进行。你传入的任何值都会被原样记录在 streamer.checkpoint.reset.key.v2 下。
直接传入一个裸时间戳会被拒绝:
Illegal checkpoint key override `20250110120000000`. Valid format is either
`resumeFromInstantRequestTime:<checkpoint value>` or `resumeFromInstantCompletionTime:<checkpoint value>`.note
S3EventsHoodieIncrSource 和 GcsEventsHoodieIncrSource 不参与完成时间检查点。无论表版本如何,它们都继续接受普通的 --checkpoint 值。
caution
使用 Hudi 增量源时,在升级或降级表版本的那次运行中,不要传入 --checkpoint 或 --ignore-checkpoint。Hudi 不会冒着在错误的检查点语义下应用该覆盖值的风险,而是会终止作业,并给出一条提示你移除这些选项的消息。请先让升级完成,然后在后续的运行中重置检查点。
转换器
Hudi Streamer 支持在写入存储之前对记录进行自定义转换。这可以通过 --transformer-class 选项提供 org.apache.hudi.utilities.transform.Transformer 的实现来完成。
SQL 查询转换器
你可以传入一个 SQL 查询,使其在写入过程中执行。
--transformer-class org.apache.hudi.utilities.transform.SqlQueryBasedTransformer
--hoodie-conf hoodie.streamer.transformer.sql=SELECT a.col1, a.col3, a.col4 FROM <SRC> aSQL 文件转换器
你可以在写入过程中指定一个包含要执行的 SQL 脚本的文件。该 SQL 文件通过以下 Hudi 属性进行配置:hoodie.streamer.transformer.sql.file
查询语句应将源数据引用为名为 <SRC> 的表。
最终 SQL 语句的结果将作为写入载荷(payload)使用。
示例 Spark SQL 查询:
CACHE TABLE tmp_personal_trips AS
SELECT * FROM <SRC> WHERE trip_type='personal_trips';
SELECT * FROM tmp_personal_trips;扁平化转换器(Flattening Transformer)
该转换器可以将嵌套对象展开为扁平结构。它会对外层字段加上前缀并附加 _,以层层嵌套的方式为传入记录中的内层字段生成前缀,从而实现字段扁平化。目前尚不支持对数组进行扁平化处理。
一个示例 schema 可能如下所示,其中 name 在原始数据源中是 StructType 的嵌套字段
age as intColumn,address as stringColumn,name.first as name_first,name.last as name_last, name.middle as name_middle将配置设置为:
--transformer-class org.apache.hudi.utilities.transform.FlatteningTransformer链式转换器(Chained Transformer)
如果希望同时使用多个转换器,可以使用链式转换器(Chained Transformer)传入多个转换器,使其按顺序依次执行。
下面的示例首先将传入的记录进行扁平化处理,然后根据指定的查询执行 SQL 投影:
--transformer-class org.apache.hudi.utilities.transform.FlatteningTransformer,org.apache.hudi.utilities.transform.SqlQueryBasedTransformer
--hoodie-conf hoodie.streamer.transformer.sql=SELECT a.col1, a.col3, a.col4 FROM <SRC> aAWS DMS Transformer
该转换器专门用于 AWS DMS 数据。如果字段不存在,它会添加 Op 字段并将其值设为 I。
配置如下:
--transformer-class org.apache.hudi.utilities.transform.AWSDmsTransformer自定义 Transformer 实现
你可以通过继承此类来编写自己的自定义 transformer。
Schema Provider
默认情况下,Spark 会推断源数据的 schema,并在写入表时使用推断出的 schema。如果你需要显式定义 schema,可以使用以下 Schema Provider 之一。
Schema Registry Provider
你可以从在线注册中心获取最新的 schema。你需要传入注册中心的 URL,如有需要,还可以在 URL 中传递用户信息和凭据,例如:https://foo:bar@schemaregistry.org。凭据会被提取出来,并作为 Authorization 请求头设置在请求中。
从注册中心获取 schema 时,你可以分别指定源 schema 和目标 schema。
| 配置 | 说明 | 示例 |
|---|---|---|
| hoodie.streamer.schemaprovider.registry.url | 你正在读取的源数据的 schema | https://foo:bar@schemaregistry.org |
| hoodie.streamer.schemaprovider.registry.targetUrl | 你正在写入的目标数据的 schema | https://foo:bar@schemaregistry.org |
上述配置通过 Hudi Streamer 的 spark-submit 命令传入,如下所示:
--hoodie-conf hoodie.streamer.schemaprovider.registry.url=https://foo:bar@schemaregistry.org与 schema registry provider 配合使用的其他可选配置,例如与 SSL 存储相关的配置,以及支持对 schema registry 返回的 schema 进行自定义转换,比如通过 org.apache.hudi.utilities.schema.converter.JsonToAvroSchemaConverter 将原始 JSON schema 转换为 Avro schema。
| 配置项 | 说明 | 示例 |
|---|---|---|
| hoodie.streamer.schemaprovider.registry.schemaconverter | 要使用的自定义 schema 转换器的类名 | org.apache.hudi.utilities.schema.converter.JsonToAvroSchemaConverter |
| schema.registry.ssl.keystore.location | SSL 密钥库位置 | |
| schema.registry.ssl.keystore.password | SSL 密钥库密码 | |
| schema.registry.ssl.truststore.location | SSL 信任库位置 | |
| schema.registry.ssl.truststore.password | SSL 信任库密码 | |
| schema.registry.ssl.key.password | SSL 密钥密码 |
JDBC Schema Provider
你可以通过 JDBC 连接获取最新的 schema。
| 配置项 | 描述 | 示例 |
|---|---|---|
| hoodie.streamer.schemaprovider.source.schema.jdbc.connection.url | 要连接的 JDBC URL。你可以在 URL 中指定源特定的连接属性 | jdbc//localhost/test?user=fred&password=secret |
| hoodie.streamer.schemaprovider.source.schema.jdbc.driver.type | 用于连接到该 URL 的 JDBC 驱动程序类名 | org.h2.Driver |
| hoodie.streamer.schemaprovider.source.schema.jdbc.username | 连接所用的用户名 | fred |
| hoodie.streamer.schemaprovider.source.schema.jdbc.password | 连接所用的密码 | secret |
| hoodie.streamer.schemaprovider.source.schema.jdbc.dbtable | 要引用其 Schema 的表 | test_database.test1_table or test1_table |
| hoodie.streamer.schemaprovider.source.schema.jdbc.timeout | 驱动程序等待 Statement 对象执行完毕的秒数。为 0 表示没有限制。在写入路径中,该选项取决于 JDBC 驱动程序对 setQueryTimeout API 的实现方式,例如,h2 JDBC 驱动程序会检查每条查询的超时,而不是整个 JDBC 批次的超时。默认值为 0。 | 0 |
| hoodie.streamer.schemaprovider.source.schema.jdbc.nullable | 如果为 true,则所有列都可为空 | true |
上述配置通过以下方式传递给 Hudi Streamer 的 spark-submit 命令:--hoodie-conf hoodie.streamer.jdbcbasedschemaprovider.connection.url=jdbc:postgresql://localhost/test?user=fred&password=secret
基于文件的 Schema Provider
你可以使用 .avsc 文件来定义 schema,然后将 DFS 上的该文件指定为 schema provider。
| 配置 | 说明 | 示例 |
|---|---|---|
| hoodie.streamer.schemaprovider.source.schema.file | 你所读取的源数据的 schema | 示例 schema 文件 |
| hoodie.streamer.schemaprovider.target.schema.file | 你所写入的目标数据的 schema | 示例 schema 文件 |
Hive Schema Provider
你可以使用 Hive 表来获取源和目标的 schema。
| 配置 | 说明 |
|---|---|
| hoodie.streamer.schemaprovider.source.schema.hive.database | 用于获取源 schema 的 Hive 数据库 |
| hoodie.streamer.schemaprovider.source.schema.hive.table | 用于获取源 schema 的 Hive 表 |
| hoodie.streamer.schemaprovider.target.schema.hive.database | 用于获取目标 schema 的 Hive 数据库 |
| hoodie.streamer.schemaprovider.target.schema.hive.table | 用于获取目标 schema 的 Hive 表 |
带后处理器的 Schema Provider
SchemaProviderWithPostProcessor 会先从前面提到的某个 Schema Provider 中提取 schema,然后应用一个后处理器,在 schema 被使用之前对其进行修改。你可以通过继承该类来编写自己的后处理器:https://github.com/apache/hudi/blob/master/hudi-utilities/src/main/java/org/apache/hudi/utilities/schema/SchemaPostProcessor.java
数据源
Hudi Streamer 可以从多种多样的来源读取数据。以下是受支持的数据源列表:
分布式文件系统(DFS)
参见存储配置页面,其中列出了一些 Hudi 可以读取的 DFS 应用示例。以下是 Hudi 在 DFS 源上支持读写的文件格式。(注意:你仍然可以使用 Spark/Flink 读取器读取其他格式,然后将数据写为 Hudi 格式。)
- CSV
- AVRO
- JSON
- PARQUET
- ORC
- HUDI
对于 DFS 源,预期的行为如下:
- 对于 JSON DFS 源,你始终需要设置 schema。如果目标 Hudi 表与源文件使用相同的 schema,只需设置源 schema 即可;否则,需要同时设置源 schema 和目标 schema。
- Hudi Streamer 直接读取源基础路径(
hoodie.streamer.source.dfs.root)下的文件,并且不会将该基础路径下的分区路径作为数据集的字段使用。详细示例请参见此处。
Kafka
Hudi 可以直接从 Kafka 集群读取数据。有关如何设置具有精确一次语义、检查点(checkpointing)和插件转换的流式摄取,请参阅 Hudi Streamer 的更多细节。从 Kafka 读取数据时支持以下格式:
- AVRO:
org.apache.hudi.utilities.sources.AvroKafkaSource - JSON:
org.apache.hudi.utilities.sources.JsonKafkaSource - Proto:
org.apache.hudi.utilities.sources.ProtoKafkaSource
详情请参阅 Kafka source config。
Pulsar
Hudi Streamer 还支持通过 org.apache.hudi.utilities.sources.PulsarSource 从 Apache Pulsar 摄取数据。详情请参阅 Pulsar source config。
Amazon Kinesis
使用 JsonKinesisSource(org.apache.hudi.utilities.sources.JsonKinesisSource)可将 JSON 记录从 AWS Kinesis Data Stream 摄取到 Hudi 表中。它会并行读取每个分片(shard),在 Hudi Streamer 检查点中跟踪每个分片的进度,自动处理分片的拆分与合并,并对 Kinesis Producer Library(KPL)生成的聚合记录进行反聚合。
常用配置
所有配置键均使用前缀 hoodie.streamer.source.kinesis.。大多数用户需要的设置如下:
| 配置项 | 默认值 | 说明 |
|---|---|---|
hoodie.streamer.source.kinesis.stream.name | (必填) | Kinesis Data Streams 流名称。 |
hoodie.streamer.source.kinesis.region | (必填) | 该流所在的 AWS 区域(例如 us-east-1)。 |
hoodie.streamer.source.kinesis.starting.position | LATEST | 在尚无 checkpoint 时的起始位置。LATEST 从每个分片的最新位置开始;EARLIEST 从 TRIM_HORIZON 开始回放。 |
hoodie.streamer.source.kinesis.max.events | 5000000 | 每批次跨所有分片读取的最大记录数。可通过该配置调整批次大小。 |
hoodie.streamer.source.kinesis.partitions | 0 | 读取时使用的 Spark 分区数。0 表示每个 Kinesis 分片对应一个 Spark 分区。设置为正值可重新分区,以提升下游并行度。 |
关于凭据,该 source 使用默认的 AWS 凭据链(实例配置文件、环境变量等)。同时也可以配置自定义端点(例如 LocalStack)的认证、API 级别的限流以及重试调优 —— 完整的 hoodie.streamer.source.kinesis.* 配置项列表请参见配置参考。
Checkpoint 格式
Hudi Streamer 会将 Kinesis 的消费进度以单个 checkpoint 字符串的形式持久化在 timeline 上。每个批次都会将 checkpoint 推进到所有分片中最后一条成功读取的记录,因此失败的批次可以安全重试,既不会跳过记录,也不会重复消费。
Checkpoint 以纯文本形式编码各分片的状态:
streamName,shardId:value,shardId:value,...每个 value 是以下之一:
lastSeq—— 从打开的分片(shard)中消费的最后一个序列号。lastSeq@arrivalTime—— 同上,附带该记录的大致到达时间(epoch 毫秒),用于延迟观测。lastSeq|endSeq—— 已关闭的分片。endSeq是该分片的最终序列号,用于在分片过期且尚未被完全消费时检测数据丢失。lastSeq@arrivalTime|endSeq—— 带到达时间的已关闭分片。
示例(序列号已缩写;Kinesis 为每个分片分配一个 56 位十进制序列号):
my-stream,shardId-000000000000:49590…88898,shardId-000000000001:49590…96306你无需自己构造或解析这个字符串——它由源端自动读取和更新——但在调试、手动重置 checkpoint 或比较各分片进度时,它会很有用。
最简 spark-submit 示例
# kinesis-source.properties
hoodie.streamer.source.kinesis.stream.name=my-stream
hoodie.streamer.source.kinesis.region=us-east-1
hoodie.streamer.source.kinesis.starting.position=LATEST
# Standard Hudi write / key-gen configs
hoodie.datasource.write.recordkey.field=id
hoodie.datasource.write.partitionpath.field=event_date
hoodie.table.ordering.fields=tsspark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer \
hudi-utilities-slim-bundle-*.jar \
--props kinesis-source.properties \
--source-class org.apache.hudi.utilities.sources.JsonKinesisSource \
--table-type COPY_ON_WRITE \
--target-base-path s3://my-bucket/hudi/my-table \
--target-table my_db.my_table \
--op UPSERT \
--continuous云存储事件源
AWS S3 存储提供了事件通知服务,当 S3 存储桶中发生某些事件时会发布通知:https://docs.aws.amazon.com/AmazonS3/latest/userguide/NotificationHowTo.html AWS 会将这些事件放入 Simple Queue Service(SQS)中。Apache Hudi 提供了 S3EventsSource 和 S3EventsHoodieIncrSource,它们可以从 SQS 中读取数据,以便在新数据或变更数据一旦出现在 S3 上时立即触发/处理。更多详情请参阅 S3 source configs。
与 S3 事件源类似,Google Cloud Storage(GCS)事件源也通过 GcsEventsSource 和 GcsEventsHoodieIncrSource 提供支持。更多详情请参阅 GCS events source configs。
AWS 配置
启用 S3 事件通知 https://docs.aws.amazon.com/AmazonS3/latest/userguide/NotificationHowTo.html
下载 aws-java-sdk-sqs jar 包。
查找队列 URL 和区域以设置这些配置:
- hoodie.streamer.s3.source.queue.url=https://sqs.us-west-2.amazonaws.com/queue/url
- hoodie.streamer.s3.source.queue.region=us-west-2
使用下方示例命令所示的 Hudi Streamer 工具启动
S3EventsSource和S3EventsHoodieIncrSource:
插入来自此博客的代码示例:https://hudi.apache.org/blog/2021/08/23/s3-events-source/#configuration-and-setup
JDBC 数据源
Hudi 可以通过全量拉取表的方式从 JDBC 数据源读取数据,也可以通过带检查点(checkpointing)的方式从 JDBC 数据源进行增量读取。
| 配置项 | 说明 | 示例 |
|---|---|---|
| hoodie.streamer.jdbc.url | JDBC 连接的 URL | jdbc//localhost/test |
| hoodie.streamer.jdbc.user | 用于 JDBC 连接身份验证的用户名 | fred |
| hoodie.streamer.jdbc.password | 用于 JDBC 连接身份验证的密码 | secret |
| hoodie.streamer.jdbc.password.file | 如果你希望使用密码文件来保存连接密码 | |
| hoodie.streamer.jdbc.driver.class | JDBC 连接所使用的驱动类 | |
| hoodie.streamer.jdbc.table.name | 表名(原文未提供说明) | my_table |
| hoodie.streamer.jdbc.table.incr.column.name | 以增量模式运行时,将使用该字段来增量拉取新数据 | |
| hoodie.streamer.jdbc.incr.pull | JDBC 连接是否执行增量拉取? | |
| hoodie.streamer.jdbc.extra.options. | 用于传递通常通过 spark.read.option() 指定的额外配置 | hoodie.streamer.jdbc.extra.options.fetchSize=100 hoodie.streamer.jdbc.extra.options.upperBound=1 hoodie.streamer.jdbc.extra.options.lowerBound=100 |
| hoodie.streamer.jdbc.storage.level | 用于控制持久化级别 | 默认 = MEMORY_AND_DISK_SER |
| hoodie.streamer.jdbc.incr.fallback.to.full.fetch | 布尔值,设为 true 时,如果增量读取出现任何错误,增量拉取将回退为全量拉取 | FALSE |
SQL Sources
SQL Source org.apache.hudi.utilities.sources.SqlSource 可从任意表读取数据,主要用于回填(backfill)作业,这类作业只处理特定的分区日期。它不会将 streamer.checkpoint.key 更新为已处理的 commit,而是获取最近一次成功的 checkpoint key,并将其设为本次回填 commit 的 checkpoint,从而不会打断常规的增量处理。要获取并使用最近的增量 checkpoint,还需要为 Hudi Streamer 作业设置这个 hoodie 配置:hoodie.write.meta.key.prefixes = 'streamer.checkpoint.key'
Spark SQL 应通过以下 hoodie 配置进行设置:hoodie.streamer.source.sql.sql.query = 'select * from source_table'
使用 org.apache.hudi.utilities.sources.SqlFileBasedSource 可以将 SQL 查询写在一个文件中,从而从任意表读取数据。SQL 文件路径应通过以下 hoodie 配置进行设置:hoodie.streamer.source.sql.file = 'hdfs://xxx/source.sql'
Debezium
Hudi Streamer 可以通过摄取 Debezium 产生的变更数据捕获(CDC)事件,使 Hudi 表与上游数据库保持同步。Debezium 将每条变更作为一条 Avro 消息发布到 Kafka 主题上,并将 schema 注册到 Confluent schema registry。Debezium 源会读取该主题,将嵌套的 Debezium 变更信封(envelope)展开为普通的表字段,并把产生的插入、更新和删除应用到目标表。
每种数据库对应一个源和一个匹配的 payload 类:
| 数据库 | 源类 | Payload 类 |
|---|---|---|
| PostgreSQL | org.apache.hudi.utilities.sources.debezium.PostgresDebeziumSource | org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload |
| MySQL | org.apache.hudi.utilities.sources.debezium.MysqlDebeziumSource | org.apache.hudi.common.model.debezium.MySqlDebeziumAvroPayload |
请注意这两部分对 MySQL 的拼写方式不同:源类是 Mysql...,而 payload 类是 MySql...。
两个源都读取 Avro 格式,且都需要 schema registry,因此应将 --schemaprovider-class 设置为 org.apache.hudi.utilities.schema.SchemaRegistryProvider,并把 hoodie.streamer.schemaprovider.registry.url 指向该主题对应的 subject。由于 Hudi 在构造 Kafka 消费者之前会丢弃所有 hoodie.* 属性,schema provider 的 URL 因此永远无法传到反序列化器,所以还需要以普通属性 schema.registry.url 的形式再次提供给 Kafka 消费者。值反序列化器本身默认已是 io.confluent.kafka.serializers.KafkaAvroDeserializer,因此只有在需要覆盖它时才需要设置 hoodie.streamer.source.kafka.value.deserializer.class。
一个 PostgreSQL 表的属性文件示例:
hoodie.streamer.source.kafka.topic=postgres.public.customers
hoodie.streamer.schemaprovider.registry.url=http://localhost:8081/subjects/postgres.public.customers-value/versions/latest
bootstrap.servers=localhost:9092
auto.offset.reset=earliest
schema.registry.url=http://localhost:8081
hoodie.datasource.write.recordkey.field=id以及读取它的作业:
[hoodie]$ spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer `ls packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle-*.jar` \
--props file://${PWD}/debezium-source.properties \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.debezium.PostgresDebeziumSource \
--payload-class org.apache.hudi.common.model.debezium.PostgresDebeziumAvroPayload \
--source-ordering-field _event_lsn \
--target-base-path file:///tmp/hudi-debezium-customers \
--target-table customers \
--table-type MERGE_ON_READ \
--op UPSERT \
--continuous该示例摄取的是 PostgreSQL。如需使用 MySQL,请替换三个 Debezium 专属参数:
--source-class org.apache.hudi.utilities.sources.debezium.MysqlDebeziumSource \
--payload-class org.apache.hudi.common.model.debezium.MySqlDebeziumAvroPayload \
--source-ordering-field _event_seq \记录键必须是上游表的主键,这样后续对某一行的修改才能就地更新该行。Merge-on-Read 表适合 CDC 流产生的小而频繁的写入,不过 Copy-on-Write 同样适用。
排序。 变更事件可能乱序到达 Kafka,因此由负载决定某一行的哪个版本胜出,而不是依赖到达顺序。对于 PostgreSQL,就是 _event_lsn 中的日志序列号,负载直接从记录中读取它。对于 MySQL,源会根据 binlog 坐标 _event_bin_file 和 _event_pos 派生出 _event_seq 列,从而得到相同的全序。将实际适用的那一个列通过 --source-ordering-field 传入,因为负载在批内去重记录时比较的就是这个值。
删除。 扁平化的 _change_operation_type 列承载每个事件的 Debezium 操作类型,取值 d 表示删除。负载会将这些事件作为删除操作应用到 Hudi 表,而不是把它们写成数据行。
对于已经在流式同步的表,hoodie.debezium.override.initial.checkpoint.key 可以设置下一次运行起始的 Kafka 偏移量。它在用 Debezium 快照初始化表之后,或者 topic 被回退重置时很有用。
caution
只要 hoodie.debezium.override.initial.checkpoint.key 处于设置状态,它就会对每个批次都生效,而不仅仅是第一个批次。只要它存在,已提交的检查点就始终是这个覆盖值,因此作业会不断从同一个偏移量重启,而无法向前推进。请只在恢复流的那次运行中设置它,之后将其移除。
错误表
Hudi Streamer 支持将错误记录隔离到与目标数据表并列的另一张表中,称为“错误表”(Error table)。这便于与死信队列(DLQ)集成。错误表通过配置 hoodie.errortable.write.class 提供 org.hudi.utilities.streamer.BaseErrorTableWriter 的用户自定义子类来实现支持。详情请参阅 org.hudi.config.HoodieErrorTableConfig。
终止策略
如有需要,用户可以在 continuous 模式下配置写入后的终止策略。例如,可以配置当连续 5 次都没有从配置的源读取到新数据时优雅关闭。以下是终止策略的接口。
/**
* Post write termination strategy for Hudi Streamer in continuous mode.
*/
public interface PostWriteTerminationStrategy {
/**
* Returns whether HoodieStreamer needs to be shutdown.
* @param scheduledCompactionInstantAndWriteStatuses optional pair of scheduled compaction instant and write statuses.
* @return true if HoodieStreamer has to be shutdown. false otherwise.
*/
boolean shouldShutdown(Option<Pair<Option<String>, JavaRDD<WriteStatus>>> scheduledCompactionInstantAndWriteStatuses);
}此外,这对新表的引导初始化(bootstrapping)也可能有所帮助。与其为了处理大量输入数据而借助大型集群执行一次 bulk load 或 bulk_insert,不如以 continuous 模式启动 Hudi Streamer,并添加一个关闭策略,以便在所有数据引导完成后终止。这样,每个批次的数据量会更小,引导数据时可能就不需要大型集群了。开箱即用的方案中已经提供了具体实现:org.apache.hudi.utilities.streamer.NoNewDataTerminationStrategy。用户也可以根据需要自行实现自己的策略。
动态配置更新
当 Hudi Streamer 以 continuous 模式运行时,属性可以在每次 sync 调用之前刷新/更新。感兴趣的用户可以实现 org.apache.hudi.utilities.streamer.ConfigurationHotUpdateStrategy 来利用这一能力。
MultiTableStreamer
HoodieMultiTableStreamer 是 Hudi Streamer 的扩展,支持将多个表同时导入 Hudi 数据集。目前,它支持按顺序导入多个表,并同时兼容 COPY_ON_WRITE 和 MERGE_ON_READ 存储类型。HoodieMultiTableStreamer 的命令行参数在很大程度上与 Hudi Streamer 相同,但有一个显著区别:必须在专用配置文件夹中以独立文件的形式提供表级别的配置。为此引入了新的命令行选项:
* --config-folder
the path to the folder which contains all the table wise config files
--base-path-prefix
this is added to enable users to create all the hudi datasets for related tables under one path in FS. The datasets are then created under the path - <base_path_prefix>/<database>/<table_to_be_ingested>. However you can override the paths for every table by setting the property hoodie.streamer.ingestion.targetBasePath需要正确设置以下属性,才能使用 HoodieMultiTableStreamer 摄取数据。
hoodie.streamer.ingestion.tablesToBeIngested
comma separated names of tables to be ingested in the format <database>.<table>, for example db1.table1,db1.table2
hoodie.streamer.ingestion.targetBasePath
if you wish to ingest a particular table in a separate path, you can mention that path here
hoodie.streamer.ingestion.<database>.<table>.configFile
path to the config file in dedicated config folder which contains table overridden properties for the particular table to be ingested.表级覆盖属性的示例配置文件位于 hudi-utilities/src/test/resources/streamer-config 下。运行 HoodieMultiTableStreamer 的命令也与运行 Hudi Streamer 的方式类似。
[hoodie]$ spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieMultiTableStreamer `ls packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle-*.jar` \
--props file://${PWD}/hudi-utilities/src/test/resources/streamer-config/kafka-source.properties \
--config-folder file://tmp/hudi-ingestion-config \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.AvroKafkaSource \
--source-ordering-field impresssiontime \
--base-path-prefix file:\/\/\/tmp/hudi-streamer-op \
--target-table uber.impressions \
--op BULK_INSERT有关配置和使用 HoodieMultiTableStreamer 的详细信息,请参阅博客章节。
按需 Hive 同步(HudiHiveSyncJob)
org.apache.hudi.utilities.HudiHiveSyncJob 是一个独立的 Spark 作业,可将 Hudi 表的元数据同步到 Hive metastore,不依赖任何摄入工作流。它适用于数据回填、手动数据修正,或在直接写入之后校对 metastore 元数据。
参数
| 参数 | 是否必填 | 说明 |
|---|---|---|
--base-path / -sp | 是 | Hudi 表的基础路径。 |
--base-file-format / -bff | 否 | 基础文件格式。默认值:PARQUET。 |
--props-file-path | 否 | 包含 Hudi / Hive 同步配置的属性文件路径。 |
--hoodie-conf | 否 | 内联配置覆盖项(可重复指定)。 |
--spark-master | 否 | Spark master URL。未设置时继承自环境变量。 |
示例
spark-submit \
--packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.HudiHiveSyncJob \
hudi-utilities-slim-bundle-*.jar \
--base-path s3://my-bucket/hudi/my-table \
--base-file-format PARQUET \
--hoodie-conf hoodie.datasource.hive_sync.mode=hms \
--hoodie-conf hoodie.datasource.hive_sync.metastore.uris=thrift://hive-metastore:9083 \
--hoodie-conf hoodie.datasource.hive_sync.database=my_db \
--hoodie-conf hoodie.datasource.hive_sync.table=my_tableDataSource 写入器所接受的全部 hoodie.datasource.hive_sync.* 选项在此处同样适用。完整选项列表请参阅同步到 Hive Metastore。
评论
登录后参与评论
KnowForge