读取表
流式读取
登录后可跨设备保存划线和私人笔记登录
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。
评论
登录后参与评论
正在加载评论…
KnowForge