)
Apache Druid Kafka Lookups 深入实践借助 Kafka 主题实现维度值的实时重命名kafka-extraction-namespace 扩展详解【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druid本文围绕 Apache Druid 的kafka-extraction-namespace扩展系统讲解如何通过订阅一个 Kafka 主题主题中每条消息的 key 是旧维度值、message 是期望的新值构建 LookupExtractorFactory从而实现维度值到人可读名称的实时映射与重命名。读完本文你将掌握该扩展的完整配置参数、源码级工作原理、集群动态配置接入方式、内存与缓存注意事项以及用 Kafka 生产者控制台端到端验证重命名功能的方法。功能定位与适用场景Lookups查找表是 Druid 中一种将维度值可选地替换为新值的机制相关的通用概念可参见 docs/content/querying/lookups.md 与 Lookup DimensionSpecs。在 lookup 语境下key 指待匹配的维度值value 指其替换后的值。例如把appid-12345重命名为Super Mega Awesome App时key 为appid-12345value 为Super Mega Awesome App。kafka-extraction-namespace正是这样一种 lookup 数据源它使 Druid 能够从 Kafka 主题中持续读取名称/键值对从而支持对维度值的重命名。一个典型场景是把用户 ID 重命名为人类可读的名称如把内部user_77483渲染为United States。与 URI/JDBC 轮询型全局缓存 lookup 相比Kafka 数据源可以在消息发布后以近乎实时的速度把更新推送到 Druid 节点适合对更新及时性要求较高的场景。需要特别说明的是Lookups 整体属于实验性experimental特性官方警告其 API 可能在向后不兼容的方向上发生变化详见 docs/content/development/experimental.md。在生产环境启用前请充分评估。前置条件加载所需扩展使用 Kafka lookup 之前必须确保同时加载两个扩展druid-lookups-cached-global提供全局缓存的命名空间机制namespace cacheKafka 数据最终也要落入这套缓存druid-kafka-extraction-namespace提供本扩展的工厂实现。加载方式是在各节点的common.runtime.properties中通过druid.extensions.loadList声明核心扩展已随发行包提供无需额外下载例如druid.extensions.loadList[druid-lookups-cached-global, druid-kafka-extraction-namespace]具体加载流程参见 Loading extensions。从模块工程文件 extensions-core/kafka-extraction-namespace/pom.xml 可以看到两个值得注意的实现细节该扩展直接依赖druid-lookups-cached-globalartifactIddruid-lookups-cached-global/artifactId因为KafkaLookupExtractorFactory通过NamespaceExtractionCacheManager创建缓存它基于org.apache.kafka:kafka_2.10:0.8.2.1即 Kafka 0.8 的 Java 客户端使用 Zookeeper 消费者连接器这与文档中仅支持 zookeeper 连接器的约束一致。扩展的接入点在 KafkaExtractionNamespaceModule.java该类实现DruidModule通过SimpleModule(kafka-lookups).registerSubtypes(KafkaLookupExtractorFactory.class)将工厂注册进 Jackson从而让配置 JSON 中的type:kafka能被正确反序列化。工作原理与消息格式如果你需要让更新尽可能及时地被消费可以订阅一个 Kafka 主题该主题的key 是旧值old value、message 是期望的新值new value两者均为 UTF-8 编码并把它作为LookupExtractorFactory接入。Druid 节点会消费该主题的全部消息并持续维护一张旧值 → 新值的映射表用于查询时的维度值替换。在源码 KafkaLookupExtractorFactory.java 的start()方法中可以看到完整的消费流程将kafkaProperties全部复制到新的Properties中并校验其中不得包含group.id与auto.offset.reset校验必须提供zookeeper.connectPreconditions.checkNotNull(..., zookeeper.connect required property)通过buildConnector(...)构造kafka.javaapi.consumer.ZookeeperConsumerConnector并使用Whitelist(Pattern.quote(topic))作为主题过滤器、DEFAULT_STRING_DECODER作为 key/message 的 UTF-8 解码器创建消息流如果主题没有产生任何流或产生了超过 1 个流都会抛出异常Topic [%s] had no streams/has %d streams! expected 1即扩展按单消费者、单流设计遍历MessageAndMetadata取出key()与message()跳过 key 或 message 为 null 的坏消息随后写入缓存映射map.put(key, message)见 L221-L232若connectTimeout 0start()会以 100ms 为粒度等待首个流真正建立连接超时则抛TimeoutException(Failed to connect to kafka in sufficient time)并回滚关闭执行器与缓存。该扩展整体运行在一个独立的后台线程Execs.singleThreaded(kafka-factory-topic-%s, Thread.MIN_PRIORITY)中以最小优先级持续消费避免影响查询线程。配置详解最小配置示例以下 JSON 即文档给出的最小可用配置{ type: kafka, kafkaTopic: testTopic, kafkaProperties: {zookeeper.connect: somehost:2181/kafka} }参数表参数说明是否必填默认值kafkaTopic要读取数据的 Kafka 主题是无kafkaPropertiesKafka 消费者属性。至少必须指定zookeeper.connect仅支持 zookeeper 连接器是无connectTimeout等待建立初始连接的时间否0不等待isOneToOne映射是否为一一对应参见 Lookup DimensionSpecs否false对参数做几点补充说明kafkaTopic与kafkaProperties必填源码中通过Preconditions.checkNotNull强制要求这两个字段kafkaTopic required、kafkaProperties required缺失时工厂构造即失败。connectTimeout非负源码字段标注Min(0)且构造函数默认值为0不等待。注意它只控制初始连接不控制后续消费这与cachedNamespace的firstCacheTimeout等待缓存首次填充语义不同。关于isOneToOne的命名差异原文档参数表中写作isOneToOne但从源码看该工厂实际的 Jackson 属性名为injectiveKafkaLookupExtractorFactory.java 中JsonProperty private final boolean injective。配置时请使用injective: true/false。当底层映射是单射的key 与 value 均唯一例如国家代码→国家名时置为true可启用内部优化见下文MapLookupExtractor。受保护属性group.id与auto.offset.reset文档明确约束消费者属性group.id和auto.offset.reset不能在kafkaProperties中设置因为扩展会强制覆盖它们——group.id被设置为UUID.randomUUID().toString()实际为kafka-factory- kafkaTopic UUID.randomUUID()见 L120auto.offset.reset被强制为smallest。这不是约定俗成而是硬性校验。源码 L166-L177 中若检测到kafkaProperties携带这两个属性会直接抛出IAE异常消息分别为Cannot set kafka property [group.id]. Property is randomly generated for you. Found [...] Cannot set kafka property [auto.offset.reset]. Property will be forced to [smallest]. Found [...]随机生成group.id的意图在于每次工厂启动都以全新消费组从头smallest消费从而保证本地缓存能从主题起始位置开始完整重建实现整个流灌入本地缓存的语义同时避免多个节点共享消费组导致分区被瓜分。这一点由测试 KafkaLookupExtractorFactoryTest.java 中的testStartFailsOnGroupID、testStartFailsOnAutoOffset验证。在集群中注册与使用 Kafka Lookupkafka类型的 lookup 与其它 lookup 一样通过 Coordinator 的动态配置下发到查询节点broker / router / peon / historical静态配置已不再支持。完整 API 说明见 Lookups 动态配置其核心接口为http://COORDINATOR_IP:PORT/druid/coordinator/v1/lookups/{tier}/{id}首次使用必须先向/druid/coordinator/v1/lookupsPOST一个空对象{}以初始化配置。之后可以用如下结构注册一个名为id_renamer的 Kafka lookup放入__defaulttier{ __default: { id_renamer: { type: kafka, kafkaTopic: user-id-renames, kafkaProperties: {zookeeper.connect: somehost:2181/kafka}, connectTimeout: 0, injective: false } } }这里__default是 lookup tier层级它独立于 historical 的 tier 概念。每个查询节点可通过druid.lookup.lookupTier默认__default声明自己属于哪个 lookup tierCoordinator 每druid.manager.lookups.period默认 30 秒检查一次配置变更并将各 tier 的 lookup 推送给该 tier 内的所有节点。注册完成后在查询中通过 Lookup DimensionSpec 引用它即可完成维度值替换例如{ type: extraction, dimension: user_id, outputName: user_name, extractionFn: { type: lookup, lookup: {type: kafka, ...} } }关于 Lookup DimensionSpec 的完整用法含 query-time map lookup 与集群全局 lookup 两种形态可参见 dimensionspecs.md。缓存与内存行为Limitations文档明确列出了该功能目前的三项限制务必在容量规划时考虑整个 Kafka 流都会被灌入本地缓存。扩展没有对主题做过滤或采样topic 中出现的所有唯一 key 最终都会占据缓存空间若使用 OnHeap 缓存Kafka 流中包含大量唯一 key 时很容易撑爆 Java 堆改用 OffHeap 缓存可以缓解这一问题但可存储的数据量依然存在上限当前没有任何淘汰eviction策略。缓存只会随消息持续增长不会按 LRU 等策略回收旧条目。缓存类型由查询节点broker、peon、historical上的运行时属性druid.lookup.namespace.cache.type控制取值与行为如下详见 lookups-cached-global.md属性说明默认值druid.lookup.namespace.cache.type命名空间缓存类型取offHeap或onHeap。offHeap使用临时文件内存映射文件实现堆外存储onHeap使用标准 Java Map 结构存于堆内onHeaponHeap使用 JVM 堆内的ConcurrentMap直接影响 GC 行为与堆大小规划offHeap使用 10MB 的堆内缓冲区外加 MapDB 内存映射文件位于 Java 临时目录。若全部cachedNamespacelookup 总量超过 10MB超出部分将作为 page cache 留在内存中由操作系统按常规调优策略换入换出。在源码层面缓存的落地由NamespaceExtractionCacheManager.createCache()返回的CacheHandler提供L186-L188getCache()返回的MapString,String即被AtomicReferenceMapString,String mapRef持有查询时包装为MapLookupExtractor(map, isInjective())L357-L361。另一个值得了解的机制是缓存键cache key与缓存失效getCacheKey()L364-L385记录工厂创建提取器瞬间的事件计数doubleEventCount每条消息写入前后各计数一次如果此后没有新消息到达返回稳定可复用的缓存键一旦有新消息写入则把随机 UUID 混入缓存键强制 Druid 查询缓存对相关结果失效从而保证重命名结果即时反映 Kafka 中的最新映射。测试testCacheKeyScramblesOnNewData、testCacheKeySameOnNoChange分别验证了这两种行为。生命周期与内省KafkaLookupExtractorFactory遵循LookupExtractorFactory接口的生命周期start()幂等重复调用仅告警并返回当前状态启动失败如连接超时会关闭执行器与缓存并返回false。测试testStartStop、testStartFailsFromTimeout、testStartStopStart、testStartStartStop覆盖了正常启停、超时回滚、以及停止后不可再启动等边界close()关闭后台线程、consumerConnector并cacheHandler.close()释放缓存未启动时调用也安全返回truetestStopWithoutStartreplaces(other)当主题、消费者属性、connectTimeout、injective任一不同时判定新配置需要替换旧工厂testReplaces覆盖get()返回基于当前缓存的MapLookupExtractor若工厂尚未启动抛出Not started异常testFailsGetNotStarted。在内省方面工厂通过KafkaLookupExtractorIntrospectionHandlerKafkaLookupExtractorIntrospectionHandler.java暴露活跃状态对工厂的GET /请求若后台消费 Future 尚未完成则返回200 OK活跃消费中否则返回410 GONE。结合通用的 lookup 内省端点/druid/lookups/v1/introspect/{lookupId}见 lookups.md可以监控该 lookup 的健康状态注意与cachedNamespace不同Kafka lookup 的内省不支持/keys、/values、/version等完整映射查看。端到端测试验证 Kafka 重命名功能为了验证整套配置是否生效可以使用 Kafka 自带的消费者/生产者控制台向主题发布 key/value 对。文档给出的生产者命令如下假设 Kafka 安装在本地且 broker 监听localhost:9092./bin/kafka-console-producer.sh --property parse.keytrue --property key.separator- --broker-list localhost:9092 --topic testTopic随后在控制台中逐行输入OLD_VAL-NEW_VAL并按回车换行即可发布一条重命名消息。例如appid-12345-Super Mega Awesome App user_77483-United States每条消息以-分隔 key 与 value正对应扩展消费时key旧值、message新值的约定。发布后Druid 节点的本地缓存会在数秒内取决于消费吞吐被更新再对相关维度执行一次携带 lookup 的查询即可看到appid-12345被替换为Super Mega Awesome App。由于auto.offset.resetsmallest且每次启动都是全新消费组即使节点在主题已有大量历史消息时才启动缓存也能从主题起始位置完整重建。如果希望快速验证受保护属性与必填项的校验逻辑可以直接运行仓库中的单元测试 KafkaLookupExtractorFactoryTest.javatestStartFailsOnGroupID、testStartFailsOnAutoOffset、testStartFailsOnMissingConnect、testSerDe验证 JSON 序列化往返分别覆盖了这些约束。注意事项小结该功能依赖 Kafka 0.8 的 Zookeeper 消费者连接器见 pom.xml 中的kafka_2.10:0.8.2.1依赖配置kafkaProperties时请使用符合该版本的连接参数kafkaProperties中除受保护属性外其它 Kafka consumer 属性如zookeeper.session.timeout.ms等均可按需透传但仍建议以zookeeper.connect为主由于没有淘汰策略且全量灌入缓存务必结合druid.lookup.namespace.cache.type推荐offHeap与主题消息体量做容量评估避免长期运行后内存持续增长Lookups 为实验性特性升级 Druid 版本时需关注相关 API 的兼容性变化。【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druid创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考