DataHub Java SDK V1 实战:用 REST、Kafka 与 File Emitter 编程式推送元数据
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
DataHub 的 Java SDK V1(io.acryl:datahub-client) 为 JVM 体系提供了一组轻量级的底层元数据发射器(Emitter),支持以编程方式直接构造并推送元数据事件到 DataHub。本文以官方文档 as-a-library.md 为核心骨架,结合仓库内 datahub-client 模块 的真实源码与测试,系统讲解 REST、Kafka、File 三种 Emitter 的安装、配置、用法与底层实现,帮助你在 CI/CD 流水线、自定义编排器、Spark 血缘等推送型场景中,把元数据事件稳定地送达 DataHub。
注意:本文描述的是 Java SDK V1,它提供的是面向元数据事件的底层 Emitter API。对于新项目,官方推荐使用Java SDK V2,其提供类型安全的实体构建器(fluent API)、简化的 CRUD 操作、基于 Patch 的高效更新,以及与 DataHub 实体模型更紧密的集成。需要从 V1 迁移可参考 迁移指南。
为什么需要编程式发射元数据
在多数场景下,元数据通过 DataHub 的 Ingestion 框架(Python 采集器)定期抓取。但在某些情况下,你需要直接构造元数据事件并以编程方式发射到 DataHub。这类需求通常是“推送型”(push-based)的,典型用例包括:
- CI/CD 流水线:构建产物、部署信息、代码变更等元数据随流水线执行即时上报;
- 自定义编排器:在任务编排的关键节点主动登记数据集、任务(DataJob)与血缘信息;
- 数据管道集成:例如仓库中的 Spark 血缘集成(acryl-spark-lineage)正是使用 Java Emitter 从 Spark 作业中发射元数据事件。
io.acryl:datahub-clientJava 包提供了 REST Emitter API,可以轻松地从任何 JVM 系统发射元数据。官方 API 指南的教程中通常带有| Java |标签页,其中大量使用了 Java API SDK 的示例。
安装:声明依赖
在你的构建系统中声明对io.acryl:datahub-client的依赖即可。动手前请先在 Maven 仓库确认io.acryl:datahub-client的最新版本号,将下面示例中的__version__替换为实际版本。
Gradle
在build.gradle中添加:
implementation 'io.acryl:datahub-client:__version__'Maven
在pom.xml中添加:
<!-- https://mvnrepository.com/artifact/io.acryl/datahub-client --> <dependency> <groupId>io.acryl</groupId> <artifactId>datahub-client</artifactId> <!-- replace __version__ with the latest version number --> <version>__version__</version> </dependency>该模块在仓库中的源码位于 metadata-integration/java/datahub-client,包结构覆盖rest、kafka、file、s3四类客户端,以及v2目录下的 SDK V2 实现。
REST Emitter:直连 DataHub 元数据服务
REST Emitter 是Apache HttpClient之上的一层轻量封装,核心实现在 RestEmitter.java。它支持非阻塞地发射元数据,并负责处理元数据 Aspect 在网络上传输时的 JSON 序列化细节。
构建与配置参数
REST Emitter 采用基于 lambda 的 fluent builder 模式构建,配置文件参数大部分与 Python 侧 datahub sink 的配置项 对应:
import datahub.client.rest.RestEmitter; //... RestEmitter emitter = RestEmitter.create(b -> b .server("http://localhost:8080") //Auth token for DataHub Cloud .token(AUTH_TOKEN_IF_NEEDED) //Override default timeout of 10 seconds .timeoutSec(OVERRIDE_DEFAULT_TIMEOUT_IN_SECONDS) //Add additional headers .extraHeaders(Collections.singletonMap("Session-token", "MY_SESSION")) // Customize HttpClient's connection ttl .customizeHttpAsyncClient(c -> c.setConnectionTimeToLive(30, TimeUnit.SECONDS)) );结合 RestEmitterConfig.java 源码,各配置项的默认值如下:
| 配置项 | 默认值 | 说明 |
|---|---|---|
server | http://localhost:8080 | DataHub GMS 服务地址 |
timeoutSec | null(底层默认 10 秒) | 覆盖默认超时。源码中DEFAULT_CONNECT_TIMEOUT_SEC与DEFAULT_READ_TIMEOUT_SEC均为 10 秒,构建器在初始化时即设置了连接请求超时与响应超时;一旦显式传入timeoutSec,会以timeoutSec * 1000毫秒覆盖之(见 RestEmitter.java) |
token | null | DataHub Cloud / 鉴权场景下的 Bearer Token,构造请求时自动附加Authorization: Bearer <token>头 |
extraHeaders | 空 Map | 额外请求头,逐项写入每个请求 |
disableSslVerification | false | 置为true时使用TrustAllStrategy与NoopHostnameVerifier关闭 SSL 证书校验(见 RestEmitter.java) |
disableChunkedEncoding | false | 置为true时关闭内容压缩,以字节数组方式提交请求体 |
maxRetries/retryIntervalSec | 0/10 | 自定义DatahubHttpRequestRetryStrategy的重试次数与间隔(秒) |
asyncIngest | null | 非空时在 payload 中附加async字段,指示 GMS 异步写入 |
asyncHttpClientBuilder | 自动构建 | 通过customizeHttpAsyncClient(...)可深度定制底层 HttpClient(如设置连接 TTL) |
使用示例
发射一个MetadataChangeProposal(MCP)元数据变更提案:
import com.linkedin.dataset.DatasetProperties; import com.linkedin.events.metadata.ChangeType; import datahub.event.MetadataChangeProposalWrapper; import datahub.client.rest.RestEmitter; import datahub.client.Callback; // ... followed by // Creates the emitter with the default coordinates and settings RestEmitter emitter = RestEmitter.createWithDefaults(); MetadataChangeProposalWrapper mcpw = MetadataChangeProposalWrapper.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.my-dataset.user-table,PROD)") .upsert() .aspect(new DatasetProperties().setDescription("This is the canonical User profile dataset")) .build(); // Blocking call using Future.get() MetadataWriteResponse requestFuture = emitter.emit(mcpw, null).get(); // Non-blocking using callback emitter.emit(mcpw, new Callback() { @Override public void onCompletion(MetadataWriteResponse response) { if (response.isSuccess()) { System.out.println(String.format("Successfully emitted metadata event for %s", mcpw.getEntityUrn())); } else { // Get the underlying http response HttpResponse httpResponse = (HttpResponse) response.getUnderlyingResponse(); System.out.println(String.format("Failed to emit metadata event for %s, aspect: %s with status code: %d", mcpw.getEntityUrn(), mcpw.getAspectName(), httpResponse.getStatusLine().getStatusCode())); // Print the server side exception if it was captured if (response.getServerException() != null) { System.out.println(String.format("Server side exception was %s", response.getServerException())); } } } @Override public void onFailure(Throwable exception) { System.out.println( String.format("Failed to emit metadata event for %s, aspect: %s due to %s", mcpw.getEntityUrn(), mcpw.getAspectName(), exception.getMessage())); } });几点实战要点:
MetadataChangeProposalWrapper是 V1 中最常用的构建入口,.entityType()声明实体类型、.entityUrn()给出唯一资源名(URN)、.upsert()表示存在即更新、.aspect()填充具体的元数据 Aspect 对象(如DatasetProperties);emitter.emit(mcpw, null)返回Future<MetadataWriteResponse>,.get()为阻塞式等待;传入Callback则为非阻塞回调式,onCompletion与onFailure分别处理成功/失败分支;- 底层实现中,每次发射实际是向
{server}/aspects?action=ingestProposal发送一个POST请求,请求头固定包含Content-Type: application/json、X-RestLi-Protocol-Version: 2.0.0与Accept: application/json(见 RestEmitter.java)。testConnection()则通过GET {server}/config探测服务可用性(见 RestEmitter.java)。
单元测试印证
仓库自带 RestEmitterTest.java,配合 TestDataHubServer.java 以本地 Mock 服务验证了发射流程、响应映射与回调行为,可作为接入时的参考范式。
Kafka Emitter:借助消息总线解耦元数据生产
Kafka Emitter 是confluent-kafka的SerializingProducer之上的一层轻量封装,提供非阻塞接口将元数据事件发送到 DataHub。核心实现在 KafkaEmitter.java。
适用场景与重要约定
当你希望将元数据生产者与 DataHub 元数据服务的可用性解耦时使用它:Kafka 作为高可用消息总线,即使 DataHub 元数据服务因计划内或意外宕机,你依然可以向 Kafka 持续收集关键系统的元数据。当发射吞吐量比“元数据已持久化到 DataHub 后端”的确认更重要时,也应选用 Kafka Emitter。
重要约定:Kafka Emitter 使用Avro对元数据事件进行序列化后发往 Kafka。DataHub 目前期望 Kafka 上的元数据事件以 Avro 序列化,更换序列化器将导致事件无法被处理。
使用示例
import java.io.IOException; import java.util.concurrent.ExecutionException; import com.linkedin.dataset.DatasetProperties; import datahub.client.kafka.KafkaEmitter; import datahub.client.kafka.KafkaEmitterConfig; import datahub.event.MetadataChangeProposalWrapper; // ... followed by // Creates the emitter with the default coordinates and settings KafkaEmitterConfig.KafkaEmitterConfigBuilder builder = KafkaEmitterConfig.builder(); KafkaEmitterConfig config = builder.build(); KafkaEmitter emitter = new KafkaEmitter(config); //Test if topic is available if(emitter.testConnection()){ MetadataChangeProposalWrapper mcpw = MetadataChangeProposalWrapper.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.my-dataset.user-table,PROD)") .upsert() .aspect(new DatasetProperties().setDescription("This is the canonical User profile dataset")) .build(); // Blocking call using future Future<MetadataWriteResponse> requestFuture = emitter.emit(mcpw, null).get(); // Non-blocking using callback emitter.emit(mcpw, new Callback() { @Override public void onFailure(Throwable exception) { System.out.println("Failed to send with: " + exception); } @Override public void onCompletion(MetadataWriteResponse metadataWriteResponse) { if (metadataWriteResponse.isSuccess()) { RecordMetadata metadata = (RecordMetadata) metadataWriteResponse.getUnderlyingResponse(); System.out.println("Sent successfully over topic: " + metadata.topic()); } else { System.out.println("Failed to send with: " + metadataWriteResponse.getUnderlyingResponse()); } } }); } else { System.out.println("Kafka service is down."); }配置项与底层行为
结合 KafkaEmitterConfig.java 源码,配置默认值如下:
| 配置项 | 默认值 | 说明 |
|---|---|---|
bootstrap | localhost:9092 | Kafka bootstrap servers 地址 |
schemaRegistryUrl | http://localhost:8081 | Schema Registry 地址,用于 Avro 序列化 |
schemaRegistryConfig | 空 Map | 透传给 ConfluentKafkaAvroSerializer的 Schema Registry 配置 |
producerConfig | 空 Map | 追加到 Kafka Producer 的任意配置项 |
initializationRetryCount | 5 | Producer 构造总尝试次数(含首次) |
initializationRetryBackoffMs | 500 | 初始化重试初始退避毫秒数 |
initializationRetryMaxBackoffMs | 4000 | 初始化重试最大退避毫秒数 |
initializationRetryMaxTotalWaitMs | 15000 | 初始化重试总等待上限(毫秒) |
底层细节(见 KafkaEmitter.java):
- 默认发送主题为
MetadataChangeProposal_v1(常量DEFAULT_MCP_KAFKA_TOPIC),构造函数也支持传入自定义主题名; - Value 序列化器固定为
io.confluent.kafka.serializers.KafkaAvroSerializer,并注入schema.registry.url;Key 使用StringSerializer,取值为实体的 URN(见 KafkaEmitter.java); emit会先将 MCP 通过AvroSerializer转为 AvroGenericRecord再投递,Future<RecordMetadata>被映射为Future<MetadataWriteResponse>,成功时getUnderlyingResponse()为RecordMetadata(可读取topic()、offset 等信息);testConnection()使用 KafkaAdminClient列出主题,超时 5000ms,Kafka 不可达时返回false;- Producer 的创建过程封装了带退避的重试逻辑(KafkaProducerInitializationRetry.java),提高依赖服务启动窗口期的健壮性。
仓库中的 KafkaEmitterTest.java 借助 Testcontainers 拉起 Kafka 与 Schema Registry(见 containers 目录)进行端到端验证,可供集成测试参考。
File Emitter:离线落盘、事后导入
File Emitter 将元数据变更提案事件(MCP)写入一个 JSON 文件,之后再交给 Python 侧的 Metadata File source 进行摄取,与 Python 侧的 Metadata File sink 机制类似。核心实现在 FileEmitter.java。
适用场景
当产生元数据事件的系统无法直接连接 DataHub 的 REST 服务或 Kafka broker时使用本方案:先落盘生成 JSON 文件,随后通过离线方式传输该文件,再用 Metadata File source 导入 DataHub。
使用示例
import datahub.client.file.FileEmitter; import datahub.client.file.FileEmitterConfig; import datahub.event.MetadataChangeProposalWrapper; // ... followed by // Define output file co-ordinates String outputFile = "/my/path/output.json"; //Create File Emitter FileEmitter emitter = new FileEmitter(FileEmitterConfig.builder().fileName(outputFile).build()); // A couple of sample metadata events MetadataChangeProposalWrapper mcpwOne = MetadataChangeProposalWrapper.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.my-dataset.user-table,PROD)") .upsert() .aspect(new DatasetProperties().setDescription("This is the canonical User profile dataset")) .build(); MetadataChangeProposalWrapper mcpwTwo = MetadataChangeProposalWrapper.builder() .entityType("dataset") .entityUrn("urn:li:dataset:(urn:li:dataPlatform:bigquery,my-project.my-dataset.fact-orders-table,PROD)") .upsert() .aspect(new DatasetProperties().setDescription("This is the canonical Fact table for orders")) .build(); MetadataChangeProposalWrapper[] mcpws = { mcpwOne, mcpwTwo }; for (MetadataChangeProposalWrapper mcpw : mcpws) { emitter.emit(mcpw); } emitter.close(); // calling close() is important to ensure file gets closed cleanly输出格式与实现要点
- FileEmitterConfig.java 仅需一个必填项
fileName指定输出文件路径; - 从源码看,File Emitter 以美化打印(4 空格缩进)的 JSON 数组格式输出:构造时写入
[,每个事件以逗号分隔,close()时补写]并关闭文件(见 FileEmitter.java),因此务必调用close()确保文件被干净地收尾; - 每次
emit会立即返回一个“成功”的Future(isSuccess() == true),因为写文件本身无需等待远端确认;若 emitter 已关闭再调用emit,则返回失败 Future 并触发onFailure回调; testConnection()对 File Emitter 无意义,调用会抛出UnsupportedOperationException。
关于 S3、GCS 等对象存储
目前File Emitter 仅支持写入本地文件系统。如果你有兴趣为它增加 S3、GCS 等对象存储支持,欢迎向社区贡献代码。(仓库中已存在独立的 S3Emitter.java,但官方文档中的 File Emitter 定位仍是本地文件。)
其他语言支持
Emitter API 同样支持其他语言:
- Python Emitter 使用指南:Python 侧的
MetadataChangeProposalWrapper与DataHubRestEmitter/DataHubKafkaEmitter与 Java 侧 API 语义一一对应,配置项也在 datahub sink 文档 中有完整描述(如token、timeout_sec、disable_ssl_verification等),跨语言切换时配置可以平滑迁移。
如何选择:三种 Emitter 的取舍
| Emitter | 网络依赖 | 确认机制 | 典型场景 |
|---|---|---|---|
| REST Emitter | 直连 DataHub GMS | 请求级 HTTP 响应确认 | JVM 系统内联发射、需要即时反馈的推送任务 |
| Kafka Emitter | 依赖 Kafka + Schema Registry | Producer ack(消息进入 Kafka 即算成功) | 高吞吐、需要与 DataHub 服务解耦、容忍异步落库 |
| File Emitter | 无 | 本地写盘即成功 | 离线/隔离环境,先落盘后由 Metadata File source 导入 |
在动手集成前,建议先浏览官方教程中带| Java |标签的 API 示例(覆盖 Dataset、DataJob、血缘等多种实体),并结合 datahub-client 模块的测试代码 理解每种 Emitter 在真实环境中的行为边界。若你的项目刚起步且需要类型安全、Patch 更新等更现代的能力,请优先评估 Java SDK V2。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考