写入表

流式写入

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

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

Spark Streaming

你可以使用 Spark 的结构化流式处理(Structured Streaming)写入 Hudi 表。

  • Scala
  • Python
# pyspark
# prepare to stream write to new table
streamingTableName = "hudi_trips_cow_streaming"
baseStreamingPath = "file:///tmp/hudi_trips_cow_streaming"
checkpointLocation = "file:///tmp/checkpoints/hudi_trips_cow_streaming"

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

# create streaming df
df = spark.readStream
    .format("hudi")
    .load(basePath)

# write stream to new hudi table
df.writeStream.format("hudi")
    .options(**hudi_streaming_options)
    .outputMode("append")
    .option("path", baseStreamingPath)
    .option("checkpointLocation", checkpointLocation)
    .trigger(once=True)
    .start()

博客

评论

登录后参与评论

正在加载评论…