读取表

流式读取

师成师成· 更新于 2026-09-29· 阅读 5 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Spark Streaming

Structured Streaming 读取基于 Hudi 的增量查询(Incremental Query)功能,因此流式读取可以返回那些提交(commit)和基础文件尚未被清理器(cleaner)删除的数据。你可以控制提交的保留时间。

  • Scala
  • Python
# pyspark
# reload data
inserts = sc._jvm.org.apache.hudi.QuickstartUtils.convertToStringList(
    dataGen.generateInserts(10))
df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))

hudi_options = {
    'hoodie.table.name': tableName,
    'hoodie.datasource.write.recordkey.field': 'uuid',
    'hoodie.datasource.write.partitionpath.field': 'partitionpath',
    'hoodie.datasource.write.table.name': tableName,
    'hoodie.datasource.write.operation': 'upsert',
    'hoodie.table.ordering.fields': 'ts',
    'hoodie.upsert.shuffle.parallelism': 2,
    'hoodie.insert.shuffle.parallelism': 2
}

df.write.format("hudi"). \
    options(**hudi_options). \
    mode("overwrite"). \
    save(basePath)

# read stream to streaming df
df = spark.readStream \
    .format("hudi") \
    .load(basePath)

# read stream and output results to console
spark.readStream \
    .format("hudi") \
    .load(basePath) \
    .writeStream \
    .format("console") \
    .start()

info

可以在 ForeachBatch sink 中使用 Spark SQL 执行 INSERT、UPDATE、DELETE 和 MERGE INTO 操作。写入之前目标表必须已经存在。

Flink 流式读取

Flink 可以作为流式源持续消费 Hudi 表中的新提交。通过设置 read.streaming.enabled=true 启用此功能,并可选择设置 read.start-commit。

CREATE TABLE hudi_table (
  uuid VARCHAR(20) PRIMARY KEY NOT ENFORCED,
  name VARCHAR(10),
  age INT,
  ts TIMESTAMP(3),
  `partition` VARCHAR(20)
)
PARTITIONED BY (`partition`)
WITH (
  'connector' = 'hudi',
  'path' = '${path}',
  'table.type' = 'MERGE_ON_READ',
  'read.streaming.enabled' = 'true',          -- enable streaming read
  'read.start-commit' = '20210316134557',      -- start from this instant (omit for latest)
  'read.streaming.check-interval' = '60'       -- poll interval in seconds
);

SELECT * FROM hudi_table;

面向流式读取的 Source V2

从 Hudi 1.2.0 开始,基于 FLIP-27 的 Source V2 可作为流式读取的可选实现。Source V2 参与 Flink 的检查点协议,支持更细粒度的恢复,并支持分区裁剪:

WITH (
  'connector' = 'hudi',
  'path' = '${path}',
  'read.streaming.enabled' = 'true',
  'read.source-v2.enabled' = 'true'   -- enable FLIP-27 source (Hudi 1.2.0+)
)

警告

使用旧版 source 创建的 savepoint 与 Source V2 不兼容。切换时请重新启动一个全新作业。迁移详情请参阅 Flink Source V2。

有关 Flink 流式读取的完整选项列表(速率限制、提交数量限制、CDC 模式等),请参阅 使用 Flink。

评论

登录后参与评论

正在加载评论…