批量写入
Spark DataSource API
hudi-spark 模块提供了 DataSource API,用于将 Spark DataFrame 写入 Hudi 表。
可用选项如下:
HoodieWriteConfig :
TABLE_NAME
DataSourceWriteOptions :
RECORDKEY_FIELD:主键字段。记录键(record key)在每个分区内唯一标识一条记录/行。如果希望实现全局唯一性,有两种选择:可以将数据集设置为不分区,或者利用全局索引来确保记录键与分区路径无关地保持唯一。记录键可以是单个列,也可以引用多个列。应根据它是简单键还是复合键相应地设置 KEYGENERATOR_CLASS_OPT_KEY 属性。例如:简单字段使用 "col1",复合字段使用 "col1,col2,col3,etc"。嵌套字段可以使用点号表示法指定,例如 a.b.c。
默认值:"uuid"
PARTITIONPATH_FIELD:用于对表进行分区的列。若要禁用分区,请将值设为空字符串,例如 ""。通过 KEYGENERATOR_CLASS_OPT_KEY 指定是否分区。如果分区路径需要进行 URL 编码,可以设置 URL_ENCODE_PARTITIONING_OPT_KEY。如果要同步到 Hive,还需通过 HIVE_PARTITION_EXTRACTOR_CLASS_OPT_KEY 指定。
默认值:"partitionpath"
ORDERING_FIELDS:当同一批次中的两条记录具有相同的键值时,将选择排序字段(ordering field)值最大的那条记录。如果 HoodieRecordPayload(WRITE_PAYLOAD_CLASS)使用的是默认载荷 OverwriteWithLatestAvroPayload,则传入的记录总是优先于存储中的记录,而忽略此排序字段配置。
无默认值
注意:配置键 hoodie.datasource.write.precombine.field 已废弃,请改用 hoodie.table.ordering.fields。
OPERATION:要使用的写入操作。
可选值:"upsert"(默认)、"bulk_insert"、"insert"、"delete"
TABLE_TYPE:要写入的表类型。注意:表初次创建后,在使用 Spark 的 SaveMode.Append 模式写入(更新)表时,该值必须保持一致。
可选值:COW_TABLE_TYPE_OPT_VAL(默认)、MOR_TABLE_TYPE_OPT_KEY
KEYGENERATOR_CLASS_NAME:请参阅下文的键生成章节。
示例:对 DataFrame 执行 upsert,并指定必要的字段名 recordKey => _row_key、partitionPath => partition 以及 orderingField => timestamp
inputDF.write()
.format("hudi")
.options(clientOpts) //Where clientOpts is of type Map[String, String]. clientOpts can include any other options necessary.
.option("hoodie.datasource.write.recordkey.field", "_row_key")
.option("hoodie.datasource.write.partitionpath.field", "partition")
.option("hoodie.table.ordering.fields", "timestamp")
.option("hoodie.table.name", tableName)
.mode(SaveMode.Append)
.save(basePath);- Scala
- Python
- SparkSQL
生成一些新的行程数据,将其加载到 DataFrame 中,然后按照如下方式将 DataFrame 写入 Hudi 表。
// spark-shell
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Overwrite).
save(basePath)info
mode(Overwrite) 会在表已存在时覆盖并重新创建该表。你可以在 /tmp/hudi_trips_cow/<region>/<country>/<city>/ 下查看生成的数据。我们提供了记录键(schema 中的 uuid)、分区字段(region/country/city)以及合并逻辑(schema 中的 ts),以确保行程记录在每个分区内都是唯一的。更多信息请参阅 Hudi 中存储数据的建模方式;关于将数据摄入 Hudi 的方式,请参阅 写入 Hudi 表。这里我们使用的是默认的写入操作:upsert。如果你的工作负载不涉及更新,也可以执行 insert 或 bulk_insert 操作,速度可能会更快。了解更多信息,请参阅 写入操作。
请查看 https://hudi.apache.org/blog/2021/02/13/hudi-key-generators,了解各种键生成器选项,例如基于时间戳的、复杂键、自定义键、非分区键生成器等。
插入覆盖表(Insert Overwrite Table)
生成一些新的行程数据,在 Hudi 元数据层面逻辑上覆盖该表。Hudi 清理器最终会清理先前表快照的文件组。这可能比删除旧表并以 Overwrite 模式重新创建更快。
- Scala
- SparkSQL
// spark-shell
spark.
read.format("hudi").
load(basePath).
select("uuid","partitionpath").
show(10, false)
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.datasource.write.operation","insert_overwrite_table").
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)
// Should have different keys now, from query before.
spark.
read.format("hudi").
load(basePath).
select("uuid","partitionpath").
show(10, false)插入覆盖
生成一些新的行程数据,覆盖输入中出现的所有分区。对于需要一次性重新计算整个目标分区的批处理 ETL 作业(而不是增量更新目标表),此操作可能比 upsert 更快。这是因为我们可以完全绕过 upsert 写入路径中的索引、预合并和其他重新分区步骤。
- Scala
- SparkSQL
// spark-shell
spark.
read.format("hudi").
load(basePath).
select("uuid","partitionpath").
sort("partitionpath","uuid").
show(100, false)
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.
read.json(spark.sparkContext.parallelize(inserts, 2)).
filter("partitionpath = 'americas/united_states/san_francisco'")
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.datasource.write.operation","insert_overwrite").
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)
// Should have different keys now for San Francisco alone, from query before.
spark.
read.format("hudi").
load(basePath).
select("uuid","partitionpath").
sort("partitionpath","uuid").
show(100, false)删除
Hudi 支持对存储在 Hudi 表中的数据实现两种类型的删除,方式是允许用户指定不同的记录载荷(record payload)实现。更多信息请参阅 Hudi 中的删除支持。
- 软删除:保留记录键(record key),仅将其余所有字段的值置为 null。只需确保表结构中相应字段可为空,然后将这些字段设为 null 并对表执行 upsert 即可实现。请注意,软删除的数据始终会持久化在存储中,不会被移除,只是所有值都被设置为 null。因此,出于 GDPR 或其他合规性原因,如果记录键和分区路径中包含个人身份信息(PII),用户应考虑执行硬删除。
例如:
// fetch two records for soft deletes
val softDeleteDs = spark.sql("select * from hudi_trips_snapshot").limit(2)
// prepare the soft deletes by ensuring the appropriate fields are nullified
val nullifyColumns = softDeleteDs.schema.fields.
map(field => (field.name, field.dataType.typeName)).
filter(pair => (!HoodieRecord.HOODIE_META_COLUMNS.contains(pair._1)
&& !Array("ts", "uuid", "partitionpath").contains(pair._1)))
val softDeleteDf = nullifyColumns.
foldLeft(softDeleteDs.drop(HoodieRecord.HOODIE_META_COLUMNS: _*))(
(ds, col) => ds.withColumn(col._1, lit(null).cast(col._2)))
// simply upsert the table after setting these fields to null
softDeleteDf.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.datasource.write.operation", "upsert").
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)- 硬删除(Hard Deletes):一种更强力的删除方式是从表中物理移除记录的所有痕迹。可以通过以下 3 种不同的方式实现。
- 使用数据源(Datasource),将
"hoodie.datasource.write.operation"设置为"delete"。这将删除所提交 DataSet 中的所有记录。
示例,首先读入一个数据集:
val roViewDF = spark.
read.
format("org.apache.hudi").
load(basePath + "/*/*/*/*")
roViewDF.createOrReplaceTempView("hudi_ro_table")
spark.sql("select count(*) from hudi_ro_table").show() // should return 10 (number of records inserted above)
val riderValue = spark.sql("select distinct rider from hudi_ro_table").show()
// copy the value displayed to be used in next step现在编写一个查询,指定你想要删除的记录:
val df = spark.sql("select uuid, partitionPath from hudi_ro_table where rider = 'rider-213'")最后,执行这些记录的删除:
val deletes = dataGen.generateDeletes(df.collectAsList())
val df = spark.read.json(spark.sparkContext.parallelize(deletes, 2));
df.write.format("org.apache.hudi").
options(getQuickstartWriteConfigs).
option("hoodie.datasource.write.operation","delete").
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath);- 使用 DataSource 时,将
PAYLOAD_CLASS_OPT_KEY设置为"org.apache.hudi.EmptyHoodieRecordPayload"。这将删除所提交 DataSet 中的所有记录。
此示例将删除表中所有存在于 DataSet deleteDF 中的记录:
deleteDF // dataframe containing just records to be deleted
.write().format("org.apache.hudi")
.option(...) // Add HUDI options like record-key, partition-path and others as needed for your setup
// specify record_key, partition_key, ordering_fields & usual params
.option(DataSourceWriteOptions.PAYLOAD_CLASS_OPT_KEY, "org.apache.hudi.EmptyHoodieRecordPayload")- 使用 DataSource 或 Hudi Streamer,为 DataSet 添加一个名为
_hoodie_is_deleted的列。对于所有需要删除的记录,该列的值必须设为true;对于需要执行 upsert 的记录,则设为false或保持为 null。
假设原始 schema 如下:
{
"type":"record",
"name":"example_tbl",
"fields":[{
"name": "uuid",
"type": "String"
}, {
"name": "ts",
"type": "string"
}, {
"name": "partitionPath",
"type": "string"
}, {
"name": "rank",
"type": "long"
}
]}确保添加 _hoodie_is_deleted 列:
{
"type":"record",
"name":"example_tbl",
"fields":[{
"name": "uuid",
"type": "String"
}, {
"name": "ts",
"type": "string"
}, {
"name": "partitionPath",
"type": "string"
}, {
"name": "rank",
"type": "long"
}, {
"name" : "_hoodie_is_deleted",
"type" : "boolean",
"default" : false
}
]}然后,对于任何想要删除的记录,你可以将 _hoodie_is_deleted 标记为 true:
{"ts": 0.0, "uuid": "19tdb048-c93e-4532-adf9-f61ce6afe10", "rank": 1045, "partitionpath": "americas/brazil/sao_paulo", "_hoodie_is_deleted" : true}写入 VECTOR、BLOB 和 VARIANT 列
VECTOR、BLOB 和 VARIANT 列可以通过 SQL INSERT 写入,参见 SQL DML;对应的 DataFrame API 用法如下。
通过 DataFrame 写入 VECTOR
在 VECTOR 列上添加 hudi_type 元数据,以便写入器能够识别该列:
import pyarrow as pa
schema = pa.schema([
pa.field("product_id", pa.string()),
pa.field("embedding", pa.list_(pa.float32()),
metadata={b"hudi_type": b"VECTOR(768)"}),
])通过 DataFrame 写入 BLOB
BLOB 列在内部是一个 struct(参见 BLOB)。将其构建为 Spark Row:
from pyspark.sql import Row
with open("logo.png", "rb") as f:
raw_bytes = f.read()
row = Row(
asset_id="asset_001",
file_name="logo.png",
mime_type="image/png",
file_size=len(raw_bytes),
content=Row(type="INLINE", data=raw_bytes, reference=None),
)对于 PyArrow schema,请显式声明该 struct:
import pyarrow as pa
schema = pa.schema([
pa.field("asset_id", pa.string()),
pa.field("file_name", pa.string()),
pa.field("mime_type", pa.string()),
pa.field("file_size", pa.int64()),
pa.field("content", pa.struct([
pa.field("type", pa.string()),
pa.field("data", pa.binary()),
pa.field("reference", pa.struct([
pa.field("external_path", pa.string()),
pa.field("offset", pa.int64()),
pa.field("length", pa.int64()),
pa.field("managed", pa.bool_()),
])),
]), metadata={b"hudi_type": b"BLOB"}),
])通过 DataFrame 写入 VARIANT
通过 DataFrame API 写入 VARIANT 需要 Spark 4.0 及以上版本。请使用原生的 VariantType:
from pyspark.sql.types import StructType, StructField, StringType, LongType, VariantType
schema = StructType([
StructField("event_id", StringType()),
StructField("payload", VariantType()),
StructField("ts", LongType()),
])或者,声明底层结构体,并在外部字段上标注 hudi_type=VARIANT。该结构体必须恰好包含两个非空的 BinaryType 字段,且字段名分别为 metadata 和 value,否则写入器会抛出 IllegalArgumentException: Invalid variant schema structure。读取时该列会以原生 VariantType 形式往返:
from pyspark.sql.types import StructType, StructField, StringType, LongType, BinaryType, MetadataBuilder
variant_metadata = MetadataBuilder().putString("hudi_type", "VARIANT").build()
variant_struct = StructType([
StructField("metadata", BinaryType(), nullable=False),
StructField("value", BinaryType(), nullable=False),
])
schema = StructType([
StructField("event_id", StringType()),
StructField("payload", variant_struct, metadata=variant_metadata),
StructField("ts", LongType()),
])从 JSON 字符串构造 VARIANT 值最简单的方式,是通过 SQL 构建 DataFrame,然后将其写出:
df = spark.sql("""
SELECT 'evt_001' AS event_id,
parse_json('{"action": "click", "x": 120, "y": 450}') AS payload,
1000 AS ts
""")
df.write.format("hudi") \
.option("hoodie.table.name", "events") \
.option("hoodie.datasource.write.recordkey.field", "event_id") \
.option("hoodie.datasource.write.precombine.field", "ts") \
.mode("append") \
.save("/path/to/table")通过 DataFrame 使用 Lance 基础文件格式
在写入选项中设置 hoodie.table.base.file.format=lance:
(df.write
.format("hudi")
.option("hoodie.table.name", "my_ai_table")
.option("hoodie.datasource.write.recordkey.field", "id")
.option("hoodie.record.merger.impls",
"org.apache.hudi.DefaultSparkRecordMerger")
.option("hoodie.table.base.file.format", "lance")
.mode("overwrite")
.save("/path/to/my_ai_table"))完整的 Lance 行为与配置参见 存储布局 → Lance。
并发控制
以下是如何通过 Spark DataSource API 使用 optimistic_concurrency_control 的示例。
关于并发控制的更多深入细节,请参阅并发控制概念章节。
inputDF.write.format("hudi")
.options(getQuickstartWriteConfigs)
.option("hoodie.table.ordering.fields", "ts")
.option("hoodie.cleaner.policy.failed.writes", "LAZY")
.option("hoodie.write.concurrency.mode", "optimistic_concurrency_control")
.option("hoodie.write.lock.zookeeper.url", "zookeeper")
.option("hoodie.write.lock.zookeeper.port", "2181")
.option("hoodie.write.lock.zookeeper.lock_key", "test_table")
.option("hoodie.write.lock.zookeeper.base_path", "/test")
.option("hoodie.datasource.write.recordkey.field", "uuid")
.option("hoodie.datasource.write.partitionpath.field", "partitionpath")
.option("hoodie.table.name", tableName)
.mode(Overwrite)
.save(basePath)滚动携带额外元数据
滚动携带额外元数据(Rolling Extra Metadata)可以让你自动将选定的提交元数据键带到之后的每一次提交和清理实例中,而无需遍历完整的时间线。这在跨提交持久化检查点信息(如 Kafka 偏移量或 Flink 检查点)时尤其有用。
| 配置项 | 默认值 | 说明 |
|---|---|---|
hoodie.write.rolling.metadata.keys | ""(已禁用) | 以逗号分隔的额外元数据键列表,这些键会被带到每一个新提交和清理实例中。值从最近已完成的实例中读取,并写入新提交的元数据中,因此无需遍历时间线即可访问。新值会覆盖旧值。仅适用于数据表的提交和清理实例。 |
hoodie.write.rolling.metadata.timeline.lookback.commits | 10 | 查找所配置的滚动元数据键时,向前回溯已完成实例的最大数量。取值越高,容错性越强,但会带来少量性能开销。 |
示例:
inputDF.write.format("hudi")
.option("hoodie.write.rolling.metadata.keys", "kafka.offset.partition.0,kafka.offset.partition.1")
.option("hoodie.write.rolling.metadata.timeline.lookback.commits", "10")
// ... other options
.save(basePath)高级存储选项
Hudi 1.2.0 中新增了以下高级存储配置选项:
| 配置项 | 默认值 | 描述 |
|---|---|---|
hoodie.parquet.write.config.injector.class | (无) | 自定义 HoodieParquetConfigInjector 实现的全限定类名。可用于注入自定义的 Parquet 写入属性(例如,禁用字典编码、设置布隆过滤器大小),而无需修改 Hudi 源码。实现类必须实现 org.apache.hudi.io.HoodieParquetConfigInjector 接口。 |
hoodie.table.base.file.format | parquet | 表的基础文件格式,可接受 parquet、orc、hfile 或 lance。有关 Lance 的特定选项,请参阅存储布局 → Lance。 |
Java 客户端
我们可以使用原生 Java 写入 Hudi 表。关于 Java 客户端的用法,请参考此处。
评论
登录后参与评论
KnowForge