数据质量
数据质量是指数据在整体上的准确性、完整性、一致性和有效性。确保数据质量对于实现准确的分析与报告、遵守相关法规,以及维护组织对自身数据基础设施的信任都至关重要。
Hudi 提供了 预提交校验器(Pre-Commit Validators),使你在使用 Hudi Streamer 或 Spark Datasource 写入器写入数据时,能够确保数据满足特定的数据质量预期。
注意
使用 BULK_INSERT 写入操作类型时,会跳过预提交校验器。
多个类名之间使用 , 分隔。
语法:hoodie.precommit.validators=class_name1,class_name2
示例:
spark.write.format("hudi")
.option("hoodie.precommit.validators", "org.apache.hudi.client.validator.SqlQueryEqualityPreCommitValidator")今天你可以使用这些校验器中的任何一个,甚至可以灵活地扩展属于自己的校验器:
SQL 单一查询结果校验器
org.apache.hudi.client.validator.SqlQuerySingleResultPreCommitValidator
SQL 单一查询结果校验器可用于校验针对表执行的查询是否返回特定值。该校验器允许你运行一条 SQL 查询,如果其结果与预期输出不匹配,则中止此次提交。
多个查询之间可以用 ; 分隔符隔开。预期结果作为查询的一部分,用 # 与查询分隔。
语法:query1#result1;query2#result2
示例:
// In this example, we set up a validator that expects there is no row with `col` column as `null`
import org.apache.hudi.config.HoodiePreCommitValidatorConfig._
df.write.format("hudi").mode(Overwrite).
option("hoodie.table.name", tableName).
option("hoodie.precommit.validators", "org.apache.hudi.client.validator.SqlQuerySingleResultPreCommitValidator").
option("hoodie.precommit.validators.single.value.sql.queries", "select count(*) from <TABLE_NAME> where col is null#0").
save(basePath)SQL 查询相等性校验
org.apache.hudi.client.validator.SqlQueryEqualityPreCommitValidator
SQL 查询相等性校验器会在数据摄取前运行一次查询,然后在数据摄取后再次运行相同的查询,并确认两次的输出结果一致。这样你就可以校验提交前后数据行的相等性。
当你需要验证某个查询不会改变数据的特定子集时,该校验器非常有用。举几个例子:
- 校验查询前后空字段的数量是否一致
- 校验查询执行后没有产生重复记录
- 校验只进行了更新操作,没有意外插入的数据
多个查询可以用 ; 分隔符分隔。
语法:query1;query2
示例:
// In this example, we set up a validator that expects no change of null rows with the new commit
import org.apache.hudi.config.HoodiePreCommitValidatorConfig._
df.write.format("hudi").mode(Overwrite).
option("hoodie.table.name", tableName).
option("hoodie.precommit.validators", "org.apache.hudi.client.validator.SqlQueryEqualityPreCommitValidator").
option("hoodie.precommit.validators.equality.sql.queries", "select count(*) from <TABLE_NAME> where col is null").
save(basePath)SQL 查询不等式校验
org.apache.hudi.client.validator.SqlQueryInequalityPreCommitValidator
SQL 查询不等式校验器会在数据写入前执行一次查询,然后在数据写入后再执行相同的查询,并确认两次的输出结果不一致。这样你就可以确认提交前后数据行的变化。
多个查询之间可以使用 ; 分隔符分隔。
语法:query1;query2
示例:
// In this example, we set up a validator that expects a change of null rows with the new commit
import org.apache.hudi.config.HoodiePreCommitValidatorConfig._
df.write.format("hudi").mode(Overwrite).
option("hoodie.table.name", tableName).
option("hoodie.precommit.validators", "org.apache.hudi.client.validator.SqlQueryInequalityPreCommitValidator").
option("hoodie.precommit.validators.inequality.sql.queries", "select count(*) from <TABLE_NAME> where col is null").
save(basePath)扩展自定义校验器
用户也可以通过扩展抽象类 SparkPreCommitValidator 并重写该方法,来提供自己的实现。
void validateRecordsBeforeAndAfter(Dataset<Row> before,
Dataset<Row> after,
Set<String> partitionsAffected)附加的带通知的监控功能
Hudi 提供了提交通知服务,可对其进行配置,以便在写入提交时触发通知。
提交通知服务可以与提交前校验器结合使用,在提交未通过校验时发送通知。实现方式是将校验的详细信息作为自定义值传递给 HTTP 端点。
关于校验器行为的说明
Hudi 1.2.0 引入了以下行为上的改进:
SQL 查询中的元数据字段:校验器的 SQL 现在可以在查询表达式中直接引用 Hudi 元数据字段(_hoodie_record_key、_hoodie_partition_path、_hoodie_file_name、_hoodie_commit_time、_hoodie_commit_seqno)。
空写入:空的写入提交不再导致提交前校验器报错。当写入中不存在任何记录时,校验器会被优雅地跳过。
失败策略
Hudi 1.2.0 为提交前校验器引入了可配置的失败策略:
| 配置键 | 默认值 | 说明 |
|---|---|---|
hoodie.precommit.validators.failure.policy | FAIL | 如何处理校验器失败。FAIL:抛出异常并阻止提交。WARN_LOG:输出警告日志但允许提交继续(适用于软性监控)。 |
Flink 与流式偏移量校验器
自 Hudi 1.2.0 起可用。Flink 写入端现在会像 Spark 一样,遵循 hoodie.precommit.validators 这一相同的配置键。需要在 Flink 中使用的校验器必须继承与引擎无关的 org.apache.hudi.client.validator.BasePreCommitValidator(位于 hudi-common 中),该基类可独立于 Spark 提供对提交元数据和时间线信息的访问。
针对基于 Kafka 的管道,现在提供了两个内置的流式偏移量校验器:
| 验证器类 | 引擎 | 说明 |
|---|---|---|
org.apache.hudi.sink.validator.FlinkKafkaOffsetValidator | Flink | 校验本批次写入的记录数是否与 Kafka 偏移量差值一致 |
org.apache.hudi.utilities.streamer.validator.SparkKafkaOffsetValidator | Spark / HoodieStreamer | 面向 Spark 的 Kafka 摄入管道,语义相同 |
两个验证器均使用以下配置:
| 配置项 | 默认值 | 说明 |
|---|---|---|
hoodie.precommit.validators.streaming.offset.tolerance.percentage | 0.0 | 基于偏移量的记录数校验的容差百分比。取值为 0.0 表示预期记录数(由 Kafka 偏移量差值计算得出)与实际写入记录数必须完全一致。对于带有去重的 upsert 工作负载,请设置更高的容差(例如 10.0 表示 10%)。 |
hoodie.precommit.validators.failure.policy | FAIL | 参见上文的失败策略。 |
示例(Flink):
hoodie.precommit.validators=org.apache.hudi.sink.validator.FlinkKafkaOffsetValidator
hoodie.precommit.validators.streaming.offset.tolerance.percentage=5.0
hoodie.precommit.validators.failure.policy=WARN_LOG写前校验器
写前校验器(Pre-Write Validators)在 Hudi 1.2.0 中引入,它在数据写入存储之前运行,这与在数据写入之后、但提交发布到时间线之前运行的写前提交校验器(pre-commit validators)形成对比。这样可以更早地拒绝无效操作,避免不必要的 I/O。
配置项:
| 配置键 | 默认值 | 说明 |
|---|---|---|
hoodie.prewrite.validators | "" | 以逗号分隔的完全限定类名列表,这些类需实现 org.apache.hudi.client.validator.PreWriteValidator 接口。 |
要实现自定义的写前校验器,请实现 org.apache.hudi.client.validator.PreWriteValidator 接口:
public interface PreWriteValidator {
<T> void validate(
String instantTime,
WriteOperationType writeOperationType,
HoodieTableMetaClient metaClient,
HoodieWriteConfig writeConfig,
HoodieEngineContext engineContext,
Option<HoodieData<HoodieRecord<T>>> recordsOpt) throws HoodieValidationException;
}目前尚未提供内置的预写入校验器实现;该框架旨在支持用户自行扩展。与预提交校验器不同,预写入校验器能够在任何写入 I/O 发生之前访问待写入的记录。
博客
视频
评论
登录后参与评论
KnowForge