
DataHub Java SDK V1 实战用 REST、Kafka 与 File Emitter 编程式推送元数据【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubDataHub 的 Java SDK V1io.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 groupIdio.acryl/groupId artifactIddatahub-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 HttpClients connection ttl .customizeHttpAsyncClient(c - c.setConnectionTimeToLive(30, TimeUnit.SECONDS)) );结合 RestEmitterConfig.java 源码各配置项的默认值如下配置项默认值说明serverhttp://localhost:8080DataHub GMS 服务地址timeoutSecnull底层默认 10 秒覆盖默认超时。源码中DEFAULT_CONNECT_TIMEOUT_SEC与DEFAULT_READ_TIMEOUT_SEC均为 10 秒构建器在初始化时即设置了连接请求超时与响应超时一旦显式传入timeoutSec会以timeoutSec * 1000毫秒覆盖之见 RestEmitter.javatokennullDataHub Cloud / 鉴权场景下的 Bearer Token构造请求时自动附加Authorization: Bearer token头extraHeaders空 Map额外请求头逐项写入每个请求disableSslVerificationfalse置为true时使用TrustAllStrategy与NoopHostnameVerifier关闭 SSL 证书校验见 RestEmitter.javadisableChunkedEncodingfalse置为true时关闭内容压缩以字节数组方式提交请求体maxRetries/retryIntervalSec0/10自定义DatahubHttpRequestRetryStrategy的重试次数与间隔秒asyncIngestnull非空时在 payload 中附加async字段指示 GMS 异步写入asyncHttpClientBuilder自动构建通过customizeHttpAsyncClient(...)可深度定制底层 HttpClient如设置连接 TTL使用示例发射一个MetadataChangeProposalMCP元数据变更提案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 对象如DatasetPropertiesemitter.emit(mcpw, null)返回FutureMetadataWriteResponse.get()为阻塞式等待传入Callback则为非阻塞回调式onCompletion与onFailure分别处理成功/失败分支底层实现中每次发射实际是向{server}/aspects?actioningestProposal发送一个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 FutureMetadataWriteResponse 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 源码配置默认值如下配置项默认值说明bootstraplocalhost:9092Kafka bootstrap servers 地址schemaRegistryUrlhttp://localhost:8081Schema Registry 地址用于 Avro 序列化schemaRegistryConfig空 Map透传给 ConfluentKafkaAvroSerializer的 Schema Registry 配置producerConfig空 Map追加到 Kafka Producer 的任意配置项initializationRetryCount5Producer 构造总尝试次数含首次initializationRetryBackoffMs500初始化重试初始退避毫秒数initializationRetryMaxBackoffMs4000初始化重试最大退避毫秒数initializationRetryMaxTotalWaitMs15000初始化重试总等待上限毫秒底层细节见 KafkaEmitter.java默认发送主题为MetadataChangeProposal_v1常量DEFAULT_MCP_KAFKA_TOPIC构造函数也支持传入自定义主题名Value 序列化器固定为io.confluent.kafka.serializers.KafkaAvroSerializer并注入schema.registry.urlKey 使用StringSerializer取值为实体的 URN见 KafkaEmitter.javaemit会先将 MCP 通过AvroSerializer转为 AvroGenericRecord再投递FutureRecordMetadata被映射为FutureMetadataWriteResponse成功时getUnderlyingResponse()为RecordMetadata可读取topic()、offset 等信息testConnection()使用 KafkaAdminClient列出主题超时 5000msKafka 不可达时返回falseProducer 的创建过程封装了带退避的重试逻辑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会立即返回一个“成功”的FutureisSuccess() 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 RegistryProducer 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),仅供参考