并发控制
并发控制定义了不同的写入者、读取者和表服务如何协调对 Hudi 表的访问。Hudi 通过将提交以原子方式发布到时间线来保证写入的原子性,每次提交都会打上一个实例时间戳,表示该操作被视为发生的时间。与通用的文件版本控制不同,Hudi 明确区分了发起写操作的写入进程、(重)写数据/元数据以进行优化或记账的表服务,以及执行查询和读取数据的读取者。
Hudi 提供:
- 三类进程之间的快照隔离,即它们都操作在表的一致性快照上。
- 写入者之间的乐观并发控制(OCC),以提供标准的关系数据库语义。
- 写入者与表服务之间、以及不同表服务之间的多版本并发控制(MVCC)。
- 写入者之间的非阻塞并发控制(NBCC),以提供流式语义并避免写入者之间的活锁/饥饿问题。
在本节中,我们将讨论 Hudi 支持的各种并发控制,以及它们如何在单写入者和多写入者场景下支持灵活的部署模式。我们还会介绍如何借助 Hudi Streamer、Hudi DataSource、Spark Structured Streaming 和 Spark SQL 等组件,从多个写入者向 Hudi 表中摄入数据。
提示
如果只有一个进程在该表上同时执行写入和异步/内联表服务,你可以配置进程内锁提供者,以避免分布式锁带来的开销。
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.InProcessLockProvider分布式锁
Hudi 中分布式协调的前提(与许多其他分布式系统类似)是需要一个分布式锁提供者,使进程能够并发地规划、调度和执行 Hudi 时间线上的操作。锁也用于生成 TrueTime。
外部锁通常与乐观并发控制配合使用,因为它可以防止两个或多个提交尝试并发修改同一资源时产生的冲突。当一个事务尝试修改被锁定的资源时,它必须等待直到锁被释放。
在多写入场景中,锁只在特定阶段(例如在提交写入之前或调度表服务之前)对 Hudi 表进行非常短暂的持有时获取,而不是在整个作业期间一直持有。这种方式允许多个写入者同时在同一个表上工作,从而提高并发度并避免冲突。
Hudi 提供了多个锁提供者,它们需要不同的配置。完整列表请参阅锁相关配置。
基于存储的锁提供者
基于存储的锁提供者通过直接利用表存储层中的 .hoodie/ 目录来管理并发,从而无需外部锁基础设施。这消除了对 DynamoDB、ZooKeeper 或 Hive Metastore 的依赖,降低了运维复杂度和成本。
它在单个锁文件上使用条件写入,确保同一时刻只有一个写入者持有锁,并通过基于心跳的续租和自动过期机制实现容错。基础配置只需设置提供者类名,同时还提供可选参数用于调优锁的有效期和续租间隔。
将相应的云组件包添加到你的类路径中:
- 对于 S3:
hudi-aws-bundle - 对于 GCS:
hudi-gcp-bundle - 对于 Azure(
abfs://、abfss://、wasb://、wasbs://):hudi-azure-bundle
设置以下配置:
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.StorageBasedLockProvider支持 S3、GCS 以及 Azure ADLS Gen2 / Azure Blob Storage。这种云原生设计直接利用存储特性,简化了大规模云上运维。
可选调优配置:
| 配置名称 | 默认值 | 描述 |
|---|---|---|
| hoodie.write.lock.storage.validity.timeout.secs | 300(可选) | 每个新锁的有效期(秒)。提供者会持续续租其锁,直到租约延长或发生超时。 |
Config Param: STORAGE_BASED_LOCK_VALIDITY_TIMEOUT_SECSSince Version: 1.0.2
| 配置名称 | 默认值 | 描述 |
|---|---|---|
| hoodie.write.lock.storage.renew.interval.secs | 30(可选) | 两次续租尝试之间的间隔(秒)。 |
Config Param: STORAGE_BASED_LOCK_RENEW_INTERVAL_SECSSince Version: 1.0.2
Azure 基于存储的锁
认证按以下优先级顺序进行解析:
| 优先级 | 配置键 | 描述 |
|---|---|---|
| 1(最高) | hoodie.write.lock.azure.connection.string | Azure Storage 连接字符串 |
| 2 | hoodie.write.lock.azure.sas.token | SAS 令牌(Azure 不建议在生产环境中使用) |
| 3 | hoodie.write.lock.azure.managed.identity.client.id | 用户分配的托管身份的客户端 ID(ManagedIdentityCredential) |
| 4 | hoodie.write.lock.azure.client.tenant.id + .client.id + .client.secret | 通过 ClientSecretCredential 使用的服务主体——三者必须全部设置 |
| 5(最低) | (无) | DefaultAzureCredential 凭据链(系统分配的托管身份、环境变量等) |
服务主体认证的配置示例:
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.StorageBasedLockProvider
hoodie.write.lock.azure.client.tenant.id=<your-tenant-id>
hoodie.write.lock.azure.client.id=<your-app-client-id>
hoodie.write.lock.azure.client.secret=<your-client-secret>基于 ZooKeeper 的锁提供者
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider基本配置:
| 配置项名称 | 默认值 | 描述 |
|---|---|---|
| hoodie.write.lock.zookeeper.base_path | 无 (必填) | 在 ZooKeeper 上创建锁相关 ZNode 的根路径。所有并发写入者必须使用相同的路径。Config Param: ZK_BASE_PATHSince Version: 0.8.0 |
| hoodie.write.lock.zookeeper.port | 无 (必填) | ZooKeeper 端口。Config Param: ZK_PORTSince Version: 0.8.0 |
| hoodie.write.lock.zookeeper.url | 无 (必填) | ZooKeeper 连接 URL。Config Param: ZK_CONNECT_URLSince Version: 0.8.0 |
基于 HiveMetastore 的锁提供者
hoodie.write.lock.provider=org.apache.hudi.hive.transaction.lock.HiveMetastoreBasedLockProvider基础配置:
| 配置名称 | 默认值 | 描述 |
|---|---|---|
hoodie.write.lock.hivemetastore.database |
N/A (必填) | 要获取锁的 Hive 数据库。Config Param: HIVE_DATABASE_NAMESince Version: 0.8.0 |
hoodie.write.lock.hivemetastore.table |
N/A (必填) | 要获取锁的 Hive 表。Config Param: HIVE_TABLE_NAMESince Version: 0.8.0 |
Hive Metastore URI 在运行时从 Hadoop 配置中获取。
:::note
HMS 锁提供程序需要 ZooKeeper。通过 Hudi 设置 ZooKeeper 配置:
:::
hoodie.write.lock.zookeeper.url
hoodie.write.lock.zookeeper.port或者通过 Hive 配置:
hive.zookeeper.quorum
hive.zookeeper.client.port基于 DynamoDB 的锁提供器
hoodie.write.lock.provider=org.apache.hudi.aws.transaction.lock.DynamoDBBasedLockProvider基于 Amazon DynamoDB 的锁提供器支持跨集群多写。所有选项参见 基于 DynamoDB 的锁配置。基本配置:
| 配置名称 | 默认值 | 描述 |
|---|---|---|
| hoodie.write.lock.dynamodb.endpoint_url | 无 (必填) | DynamoDB 的端点 URL(便于本地测试)。 |
配置参数:DYNAMODB_ENDPOINT_URL起始版本:0.10.1
更多配置:基于 DynamoDB 的锁配置
建表:Hudi 会自动创建 hoodie.write.lock.dynamodb.table 指定的 DynamoDB 表。若使用已有表,请确保存在 key 属性(作为分区键)。hoodie.write.lock.dynamodb.partition_key 是为分区键写入的值(未设置时根据表名推断),以确保多个写入者使用同一把锁。
凭证属性(若不使用默认提供器链):
hoodie.aws.access.key
hoodie.aws.secret.key
hoodie.aws.session.token若未设置,Hudi 将回退到 DefaultAWSCredentialsProviderChain。
所需的 IAM 权限:
{
"Sid": "DynamoDBLocksTable",
"Effect": "Allow",
"Action": [
"dynamodb:CreateTable",
"dynamodb:DeleteItem",
"dynamodb:DescribeTable",
"dynamodb:GetItem",
"dynamodb:PutItem",
"dynamodb:Scan",
"dynamodb:UpdateItem"
],
"Resource": "arn:${Partition}:dynamodb:${Region}:${Account}:table/${TableName}"
}TableName:与hoodie.write.lock.dynamodb.partition_key相同Region:与hoodie.write.lock.dynamodb.region相同
从 v0.10.x 起,AWS SDK 依赖不再随 Hudi 一起打包,需要自行添加到类路径中。请添加以下 Maven 依赖(安装时请确认最新版本):
com.amazonaws:dynamodb-lock-client
com.amazonaws:aws-java-sdk-dynamodb
com.amazonaws:aws-java-sdk-core基于 DynamoDB 的隐式分区键锁提供器
hoodie.write.lock.provider=org.apache.hudi.aws.transaction.lock.DynamoDBBasedImplicitPartitionKeyLockProvider该变体的行为与上述基于 DynamoDB 的锁提供者基本一致,唯一的区别在于确定 DynamoDB 分区键的方式。它不读取 hoodie.write.lock.dynamodb.partition_key,而是根据表的基础路径推导出键:取该路径的 64 位 xxHash,并将 s3a:// 规范化为 s3://,从而让通过任一 scheme 访问同一表的写入者获取到同一把锁。
当多个表共享同一个锁表时,建议使用该变体。标准提供者的分区键取自 hoodie.write.lock.dynamodb.partition_key,而你很少会设置它:当该配置缺失时,Hudi 会用表名来填充。因此,两个恰好同名但位于不同数据库或不同路径下的表会解析到同一把锁,从而对彼此从不涉及相同数据的写入者进行不必要的串行化。而根据基础路径推导键可以保证每个表的键唯一,无需针对每个表进行配置。
其他方面均保持不变:锁表、区域、计费模式、端点 URL、凭据和 IAM 权限均按上文所述方式工作,并且不会读取 hoodie.write.lock.dynamodb.partition_key。
note
由于行是基于哈希而非表名进行键控的,因此 DynamoDB 锁表中的条目无法一眼识别。该提供者在获取锁时会将基础路径连同推导出的键一并记录到日志中,这就是将两者对应起来的方法。
基于 FileSystem 的锁提供者
基于 FileSystem 的锁提供者利用底层文件系统的原子创建/删除操作,支持跨不同作业/应用的多个写入者。
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.FileSystemBasedLockProvider使用基于 FileSystem 的锁提供者时,锁文件默认存储在 hoodie.base.path + /.hoodie/lock 下。你也可以通过指定 hoodie.write.lock.filesystem.path 来使用自定义目录存储锁文件。
如果任务崩溃导致锁无法释放,你可以将 hoodie.write.lock.filesystem.expire(默认情况下锁永不过期)设置为所需的过期时间(单位为分钟)。在这种情况下,你也可以手动删除锁文件。
warning
基于 FileSystem 的锁提供者不适用于生产环境,且不支持 S3、GCS 等云存储。
简单单写入者 + 表服务
数据湖仓管道大多以单写入者为主,表上最常见的分布式协调需求来自表管理。例如,一个 Apache Flink 作业持续向表中快速写入数据,同时需要定期进行文件大小管理或清理。Hudi 的存储引擎和平台工具为这类常见场景提供了大量支持。
内联表服务
这是最简单的并发形式,即写入过程中完全不存在并发。在这种模式下,Hudi 通过开箱即用地支持这些表服务,并在每次写入表后以内联方式运行,从而消除了对并发控制的需求并最大化吞吐量。执行计划是幂等的,会持久化到时间线上,并能在故障后自动恢复。对于大多数简单用例,只需写入即可获得一个管理良好、无需并发控制的表。
该模型中不存在实际的并发写入。MVCC 被用于在摄取写入者与多个读取者之间、以及多个表服务写入者与读取者之间提供快照隔离保证。无论是来自摄取还是来自表服务的表写入,都会产生带版本的数据,只有在写入提交后才对读取者可见。在此之前,读取者只能访问数据的上一个版本。
单个写入者搭配所有表服务(如清理、聚类、压缩等)可以配置为内联执行(例如 Hudi Streamer 的 sync-once 模式以及采用默认配置的 Spark Datasource),无需任何额外配置。
异步表服务
Hudi 提供了以异步方式运行表服务的选项,其中大部分繁重工作(例如压缩服务实际重写列式数据)以异步方式完成。在这种模式下,异步部署消除了任何重复的浪费性重试,并在单个写入者消费表的写入数据、无需被此类表服务阻塞的同时,利用聚类技术对表进行优化。该模型避免了使用锁提供者来控制并发的需求,也避免了单独编排和监控离线表服务作业的需求。
单个写入器与异步表服务在同一进程中运行。例如,你可以让 Hudi Streamer 以连续模式运行,使用异步压缩(compaction)向 MOR 表写入数据;你可以使用 Spark Streaming(其中压缩默认是异步的),也可以使用 Flink 流式作业或自建作业,并在同一个写入器内部启用异步表服务。
在这种模型下,Hudi 利用 MVCC 来支持任意数量的表服务作业并发运行,且不会产生并发冲突。这是通过确保 Hudi 的摄取写入器与异步表服务之间相互协调,从而避免冲突与竞态条件来实现的。上述模型 A 中描述的单写入器保证在该模型下同样可以达成。采用该模型,用户无需启动不同的 Spark 作业并管理它们之间的编排。对于大规模部署,该模型可以在让表服务运行的同时不阻塞写入器,从而显著减轻运维负担。
单写入器保证
在该模型下,对写入操作的预期保证如下:
- UPSERT 保证:目标表绝不会出现重复数据。
- INSERT 保证:如果启用去重配置
hoodie.datasource.write.insert.drop.duplicates与hoodie.combine.before.insert,目标表绝不会出现重复数据。 - BULK_INSERT 保证:如果启用去重配置
hoodie.datasource.write.insert.drop.duplicates与hoodie.combine.before.insert,目标表绝不会出现重复数据。 - 增量查询保证:数据消费与检查点绝不会乱序。
完全多写入器 + 异步表服务
Hudi 支持并发模式 NON_BLOCKING_CONCURRENCY_CONTROL,与 OCC 不同,多个写入器可以在非阻塞冲突解决机制下操作同一张表。写入器可以向同一个文件组写入数据,冲突由查询读取器和压缩器自动解决。你可以在该章节中了解更多相关信息。
并非总能将对表的所有写入操作(如 UPSERT、INSERT 或 DELETE)串行化到同一个写入进程中,因此可能需要多写入能力。在多写入场景下,不同的分布式进程在并行或重叠的时间窗口内运行,向同一张表写入数据。在这种情况下,必须使用外部锁机制来安全地协调并发访问。以下几种不同的场景都属于多写入场景:
- 向同一张表进行多写入:例如,两个 Spark Datasource writer 分别处理来自同一 Kafka topic 的不同分区集合。
- 向同一张表进行多写入,其中某个 writer 带有异步表服务:例如,一个 Hudi Streamer 使用异步压缩进行常规摄取,同时一个 Spark Datasource writer 用于回填。
- 一个摄取 writer 加上一个独立于摄取 writer 之外的压缩(HoodieCompactor)或聚类(HoodieClusteringJob)任务:由于它们不在同一进程中运行,因此也被视为多写入。
Hudi 的并发模型会智能地区分"向表实际写入数据"和"管理或优化表的表服务"。Hudi 提供类似的跨多个 writer 的乐观并发控制,但只要表服务与某个 writer 运行在同一进程中,表服务仍然可以完全无锁且异步地执行。对于多写入场景,Hudi 采用文件级别的乐观并发控制(OCC)。例如,当两个 writer 写入互不重叠的文件时,两次写入都会成功。但是,当不同 writer 的写入发生重叠(涉及同一组文件)时,只有其中一个会成功。请注意,该功能目前仍处于实验阶段,需要外部锁提供方在写入过程中的临界区短暂获取锁。关于锁提供方,详见下文。
多写入保证
在多个 writer 使用 OCC 的情况下,可预期的写入保证如下:
- UPSERT 保证:目标表绝不会出现重复数据。
- INSERT 保证:即使启用了去重,目标表可能会出现重复数据。
- BULK_INSERT 保证:即使启用了去重,目标表可能会出现重复数据。
- INCREMENTAL PULL 保证:数据消费和检查点绝不会乱序。如果存在未完成的提交(由多写入导致),增量查询不会暴露排在这些未完成提交之后的已完成提交。
非阻塞并发控制
NON_BLOCKING_CONCURRENCY_CONTROL 提供与上述 OCC 场景相同的保证,但无需显式锁来串行化写入。只有在将提交元数据写入 Hudi timeline 时才需要加锁。提交的完成时间反映了串行化顺序,文件切片(file slicing)基于完成时间进行。多个 writer 可以通过非阻塞冲突解决方式操作同一张表。writer 可以写入同一个文件组,冲突由查询读取方和压缩任务自动解决。该机制同时适用于压缩和摄取场景,可以参考 Flink writers 中的示例。
注意
NON_BLOCKING_CONCURRENCY_CONTROL 目前仅适用于使用简单桶索引(simple bucket index)和分区级桶索引(partition-level bucket index)的 MOR 表。
注意
NON_BLOCKING_CONCURRENCY_CONTROL 目前尚未支持摄取写入器与表服务写入器之间的聚类操作。聚类请使用 OPTIMISTIC_CONCURRENCY_CONTROL。
早期冲突检测
使用 OCC 的多写入方式允许多个写入器在不存在待写入重叠数据文件的情况下,向 Hudi 表并发写入并原子化提交,从而保证数据的一致性、完整性与正确性。在 0.13.0 版本之前,正如 OCC(乐观并发控制)这一名称所暗示的那样,每个写入器都会乐观地推进摄取流程,直到最后、即将提交之前才会进入冲突解决流程,以推断是否存在重叠写入,并在必要时中止其中一个。但这会导致大量的计算资源浪费,因为被中止的提交必须从头重试。从 0.13.0 开始,Hudi 借助 hudi 中的标记(marker)机制引入了早期冲突推断,能够在写入生命周期中尽早推断冲突并提前中止,而不是等到最后才处理。对于大规模部署而言,如果可能存在重叠的并发写入器,这一机制可以避免浪费大量计算资源。
为了改进并发控制,0.13.0 版本引入了一项新特性——OCC 中的早期冲突检测,它利用 Hudi 的标记机制在数据写入阶段检测冲突,一旦检测到冲突便提前中止写入。得益于早期冲突检测,Hudi 现在可以更早地停止产生冲突的写入器,并释放聚类所需的计算资源,从而提升资源利用率。
该特性默认关闭。如需启用,用户在使用 OCC 进行并发控制时需要将 hoodie.write.concurrency.early.conflict.detection.enable 设置为 true(所有相关配置请参阅配置页面)。
note
OCC 中的早期冲突检测是一项实验性特性
启用多写入
要开启乐观并发控制以实现多写入,需要正确设置以下属性。
hoodie.write.concurrency.mode=optimistic_concurrency_control
hoodie.write.lock.provider=<lock-provider-classname>
hoodie.cleaner.policy.failed.writes=LAZY| Config Name | Default | Description |
|---|---|---|
hoodie.write.concurrency.mode |
SINGLE_WRITER(可选) |
写操作的并发模式。可取值如下: - SINGLE_WRITER:表上只有一个活跃写入者,可最大化吞吐量。- OPTIMISTIC_CONCURRENCY_CONTROL:多个写入者可以同时操作该表,通过锁进行延迟冲突解决。也就是说,当多个写入者写入同一个文件组时,最终只有一个写入者成功。- NON_BLOCKING_CONCURRENCY_CONTROL:多个写入者可以同时操作该表,且冲突的解决是非阻塞的。写入者可以写入同一个文件组,冲突由查询读取端和 compactor 自动解决。Config Param: WRITE_CONCURRENCY_MODE |
hoodie.write.lock.provider |
无 (必填) | 锁提供者的类名,用户可以提供自己的 LockProvider 实现,该实现必须是 org.apache.hudi.common.lock.LockProvider 的子类。Config Param: LOCK_PROVIDER_CLASS_NAMESince Version: 0.8.0 |
hoodie.cleaner.policy.failed.writes |
EAGER(可选) |
org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy:控制如何清理失败写入的策略。Hudi 会删除失败写入所产生的文件以回收空间。EAGER(默认):每次写操作完成后立即内联清理失败的写入。LAZY:在心跳超时后,由清理服务运行时延迟清理失败的写入。启用多写入者时必须使用该策略。NEVER:从不清理失败的写入。Config Param: FAILED_WRITES_CLEANER_POLICY |
通过 Hudi Streamer 实现多写入
Hudi Streamer(属于 hudi-utilities-slim-bundle 的一部分)支持从 DFS、Kafka 等不同数据源进行数据摄入,并具备以下能力。
通过 Hudi Streamer 使用 optimistic_concurrency_control,需要将上述配置添加到可以传递给作业的 properties 文件中。例如,在下面的示例中,将配置添加到 kafka-source.properties 文件并传递给 Hudi Streamer,即可启用乐观并发控制。随后可以按如下方式触发 Hudi Streamer 作业:
[hoodie]$ spark-submit \
--jars "packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar,packaging/hudi-spark-bundle/target/hudi-spark3.5-bundle_2.12-1.2.0.jar" \
--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 tableName \
--op BULK_INSERT通过 Spark DataSource 写入器实现多写
hudi-spark 模块提供 DataSource API,用于将 Spark DataFrame 写入(或读取)Hudi 表。
以下示例展示了如何通过 Spark DataSource 使用 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.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)禁用多写入
移除以下曾用于启用多写入器的设置,或将其恢复为默认值。
hoodie.write.concurrency.mode=single_writer
hoodie.cleaner.policy.failed.writes=EAGEROCC 最佳实践
要并发写入 Hudi 表,需要使用前文提到的某一种锁提供者来获取锁。出于多种原因,你可能需要配置重试机制,以便应用能够成功获取锁。
- 网络连接问题或服务器负载过重,导致获取锁的时间变长,从而引发超时;
- 大量并发作业同时向同一张 Hudi 表写入数据,可能在获取锁时产生竞争,导致超时;
- 在某些冲突解决场景中,Hudi 的提交操作可能需要十几秒的时间,而在此期间锁一直处于持有状态,这会导致其他等待获取锁的作业超时。
请为原生锁提供者客户端设置正确的重试参数。:::note 请注意,有时这些配置是在服务器端一次性设置的,所有客户端都会继承相同的配置。在启用乐观并发控制之前,请先检查你的配置。:::
hoodie.write.lock.wait_time_ms
hoodie.write.lock.num_retries为 Zookeeper 和 HiveMetastore 设置正确的 Hudi 客户端重试次数。当无法修改原生客户端的重试设置时,此配置非常有用。请注意,这些重试会在你已设置的任何原生客户端重试之外额外发生。
hoodie.write.lock.client.wait_time_ms
hoodie.write.lock.client.num_retries这些值的设置需要根据具体情况而定;对于一般场景,已提供了一些默认值。
写前清理策略
运行多写入器管道时,如果某个写入器在清理周期运行前崩溃,失败的写入就会不断累积在存储中。Hudi 1.2.0 引入了 hoodie.prewrite.cleaner.policy,以便在写入启动时主动处理这一问题:
| 配置项 | 默认值 | 描述 |
|---|---|---|
hoodie.prewrite.cleaner.policy | NONE | 在开始新的摄取写入提交之前应用的策略。NONE:不执行写前操作(默认)。CLEAN:强制执行一次表清理调用(同时会回滚失败的写入)。ROLLBACK_FAILED_WRITES:仅回滚失败的写入,不执行完整的清理。 |
当某个写入器总是未能完成 CLEAN 就崩溃时,此配置非常有用。完整的清理配置列表请参见清理。
锁审计日志与诊断
基于存储的锁提供器支持对锁操作进行可选的审计日志记录。启用后,会在表基础路径下写入一个 .hoodie/lock/audit_enabled.json 标记文件,并记录锁的获取与释放事件,便于事后调试。
对于基于 ZooKeeper 的锁,ZK 锁节点现在会存储持有该锁的写入方的 Spark 应用程序 ID,从而更方便地在集群 UI 中将锁持有者与正在运行的 Spark 作业关联起来。
注意事项
如果您使用的是 WriteClient API,请注意,对表的多次写入必须由 2 个不同的写入客户端实例 发起。不建议使用同一个写入客户端实例来执行多写入操作。
博客
视频
评论
登录后参与评论
KnowForge