数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载Delta Lake 是构建 Lakehouse 架构的开放存储框架而 Apache Druid 是高性能实时分析数据库两者结合可以把存储在 Delta 表中的数据直接摄入 Druid 进行实时分析。本文以仓库中的druid-deltalake-extensions扩展为核心讲解其基于 Delta Kernel 的底层工作原理、安装加载方式、delta输入源的完整配置方法以及 8 种 Delta 过滤器的实战用法读完后你可以直接在本地 Druid 集群中把 Delta Lake 表的最新快照摄入为 Druid 数据源。扩展概述为什么需要 Delta Lake 连接器Delta Lake 提供了事务性、可伸缩的数据湖能力支持 Spark、Flink 等多样化的计算引擎在同一份数据上工作。但 Druid 并不能直接读取 Delta 表——Delta 表中的数据以版本化 Parquet 文件的形式存储并带有事务日志_delta_log普通 Parquet 输入源无法感知其协议。为此Druid 官方提供了社区扩展 druid-deltalake-extensions其中实现了DeltaInputSource即type: delta输入源。它的作用正如 delta-lake.md 所述Delta Lake is an open source storage framework that enables building a Lakehouse architecture with various compute engines. DeltaLakeInputSource lets you ingest data stored in a Delta Lake table into Apache Druid.要使用该扩展需要将druid-deltalake-extensions添加到 Druid 的已加载扩展列表中具体加载方式参见 Loading extensions社区扩展的加载说明见同页的 Loading community extensions 小节。工作原理从最新快照到 Druid InputRow 的数据链路该扩展并没有重新实现 Delta 协议而是直接基于 Delta Kernel API从 Delta Lake 3.0.0 引入的官方内核抽象与 Delta 表交互。从 DeltaInputSource.java 的源码可以梳理出完整的摄入流程定位表通过Table.forPath(engine, tablePath)打开指定路径的 Delta 表引擎由DefaultEngine.create(conf)创建内部基于 HadoopConfiguration。取最新快照table.getLatestSnapshot(engine)获取当前最新快照及其完整 Schema。值得注意的是代码在调用该 API 前会把上下文类加载器临时切换为LogStore的类加载器这是针对 Delta Kernel 3.2.0 在实例化LogStore时的已知问题对应 delta-io/delta 的 issue 3299所做的 workaround详见源码注释。列剪枝优化pruneSchema()根据InputRowSchema的ColumnsFilter从快照 Schema 中筛选需要的列构造物理读取 Schema从而在扫描阶段就只读必要列。构建 Scan 并应用过滤通过ScanBuilder把用户配置的 Delta 过滤器翻译为 Delta Kernel 的PredicatescanBuilder.withFilter(...)并配合裁剪后的读取 Schema 构建Scan。枚举数据文件scan.getScanFiles(engine)返回当前快照中需要读取的扫描文件列表每个文件对应一个DeltaSplit其state字段保存快照状态的 JSON 序列化结果files字段保存扫描文件列表见 DeltaSplit.java。读取 Parquet 数据对每个扫描文件engine.getParquetHandler().readParquetFiles(...)按物理读取 Schema和剩余谓词读取 Parquet再经Scan.transformPhysicalData转换为列式批数据。转换为 Druid 行DeltaInputSourceReader.java 中的DeltaInputSourceIterator逐批消费列式数据每个 KernelRow被包装为 DeltaInputRow.java最终通过MapInputRowParser解析成 Druid 的InputRow。文档对这一步的概括是Delta 输入源读取配置的 Delta 表并基于可选的 Delta 过滤器提取该表最新快照中的底层 Delta 文件——这些 Delta Lake 文件本身就是带版本信息的 Parquet 文件因此DeltaInputSource.needsFormat()直接返回false即无需也不允许再额外指定输入格式格式固定为 Parquet。可拆分性并行摄入的基础DeltaInputSource实现了SplittableInputSourceDeltaSplitcreateSplits()会把最新快照的每个扫描文件封装成一个独立的InputSplitestimateNumSplits()返回分片数量withSplit()则为每个分片构造只含单个 split 的新输入源。这套机制让index_parallel任务可以把不同 Delta 文件分发到不同任务并行处理是批式摄入吞吐的关键。版本支持根据 delta-lake.md 的 Version support 小节该扩展使用Delta Lake 3.0.0 引入的 Delta Kernel其兼容Apache Spark 3.5.x更旧的 Delta Lake 版本不受支持如需使用本扩展请升级到 Delta Lake 3.0.x 或更高版本。从仓库当前状态看pom.xml 中声明的delta-kernel.version为3.2.0依赖包括delta-kernel-api、delta-kernel-defaults与delta-storage三个内核模块。同时DeltaInputSource.java 的 Javadoc 明确指出目前 Delta Kernel 的 Table API 只支持读取最新快照。安装与加载扩展与大多数 Druid 社区扩展一样druid-deltalake-extensions通过pull-deps工具下载。在替换VERSION为期望的 Druid 版本后执行如下命令命令来自原文档保持原样java \ -cp lib/* \ -Ddruid.extensions.directoryextensions \ -Ddruid.extensions.hadoopDependenciesDirhadoop-dependencies \ org.apache.druid.cli.Main tools pull-deps \ --no-default-hadoop \ -c org.apache.druid.extensions.contrib:druid-deltalake-extensions:VERSION要点说明-Ddruid.extensions.directory指定扩展安装目录默认为extensions下载的 jar 会被放到这里-Ddruid.extensions.hadoopDependenciesDir指定 Hadoop 依赖目录--no-default-hadoop表示不拉取默认 Hadoop 依赖-c后的坐标由 groupIdorg.apache.druid.extensions.contrib、artifactIddruid-deltalake-extensions和版本号组成。下载完成后还需把该扩展加入druid.extensions.loadList配置参见 Loading extensions并重启相关 Druid 服务使其生效。若使用包含全部社区扩展的发行包该扩展已随包分发只需确认其在加载列表中。使用 Delta 输入源启用扩展后即可在批式摄入任务index_parallel的ioConfig.inputSource中使用type: delta。核心属性如下表源自 input-sources.md属性描述是否必填type固定为delta是tablePathDelta 表所在位置本地路径或对象存储路径是filter用于在快照内过滤数据文件的 JSON 对象否示例一读取整个快照以下 spec 读取/delta-table/foo表中的全部记录... ioConfig: { type: index_parallel, inputSource: { type: delta, tablePath: /delta-table/foo }, }在源码层面tablePath为空时会抛出InvalidInputtablePath cannot be null.且该路径直接传给 Delta Kernel 的Table.forPath()因此可以是本地文件系统路径也可以是 Delta Kernel 支持的对象存储路径。示例二带过滤器的读取以下 spec 只读取name Employee4 and age 30的数据... ioConfig: { type: index_parallel, inputSource: { type: delta, tablePath: /delta-table/foo, filter: { type: and, filters: [ { type: , column: name, value: Employee4 }, { type: , column: age, value: 30 } ] } }, }摄入任务的其余部分dataSchema、tuningConfig等与普通批式任务一致可参考 native-batch.md 中index_parallel的完整 spec 结构。Delta 过滤器详解过滤器的作用是在快照层面剪除不需要的数据文件从而减少 Druid 需要摄入的文件数量。输入源共提供 8 种过滤器and、or、not、、、、、。从 DeltaFilter.java 可以看到该接口通过 Jackson 注解注册了全部子类型type字段即 JSON 中的过滤器名每个过滤器最终通过getFilterPredicate(snapshotSchema)翻译成 Delta Kernel 的Predicate表达式树。各过滤器参数and过滤器逻辑与两个条件都必须为真属性描述是否必填type固定为and是filtersDelta 过滤器谓词列表要求恰好两个过滤器是or过滤器逻辑或满足其一即可属性描述是否必填type固定为or是filtersDelta 过滤器谓词列表要求恰好两个过滤器是not过滤器逻辑非属性描述是否必填type固定为not是filter被取反的 Delta 过滤器要求恰好一个是比较类过滤器、、、、参数一致属性描述是否必填type分别固定为、、、、是column应用过滤器的表列名是value过滤器使用的值是过滤的语义与保证需要特别注意过滤的语义边界这是该扩展最重要的使用前提原文档与源码 Javadoc 均明确说明对分区列过滤保证生效。当过滤器作用于分区表的分区列时Delta Kernel 可以精确剪枝只读取匹配分区的文件对非分区列过滤best-effort尽力而为。Delta Kernel 只依赖建表时收集的统计信息进行剪枝因此 Druid 连接器可能摄入不符合过滤条件的数据。若要确保 Delta Kernel 能剪除不必要的列值请只在分区列上使用过滤器。过滤器实现细节类型推断DeltaFilterUtils.java 的dataTypeToLiteral()会根据快照 Schema 中列的数据类型把字符串值转换为对应的 Delta 字面量支持String、Integer、Short、Long、Float、Double、Date若列不存在或类型不支持会抛出InvalidInput。数值列若传入非数字值会提示value must be a number。组合限制从 DeltaAndFilter.java 的源码看and/or目前只允许恰好两个谓词多余或不足都会抛出InvalidInput源码注释提到未来可以通过递归展平支持更复杂的表达式树。not则要求恰好一个谓词。翻译方式以为例DeltaEqualsFilter会构造new Predicate(, [Column(column), literal])and直接构造 Kernel 的And(left, right)谓词。数据类型映射Delta Kernel 的列式Row需要转换为 Druid 的InputRow。从 DeltaInputRow.java 的getValue()可以看到支持的 Delta 类型及转换规则Delta Kernel 类型转换结果BooleanTypebooleanByteType/ShortType/IntegerType对应整数类型DateType由DeltaTimeUtils.getSecondsFromDate(...)转换为 epoch 秒LongTypelongTimestampType由DeltaTimeUtils.getMillisFromTimestamp(...)转换为 epoch 毫秒FloatType/DoubleType对应浮点类型StringType字符串BinaryType字节数组按字符转换后以字符串形式返回DecimalType以decimal.longValue()转换为long其他类型抛出InvalidInputUnsupported data type其中Date/Timestamp的时间换算逻辑集中在 DeltaTimeUtils.java这决定了 Delta 表中的时间列进入 Druid 后的数值语义设计timestampSpec时应与其保持一致。行转换完成后DeltaInputRow委托给MapInputRowParser完成 Druid 维度/指标/时间戳的解析因此下游的timestampSpec、dimensionsSpec、metricsSpec用法与其他输入源完全一致。已知限制综合 delta-lake.md 的 Known limitations 小节与源码注释使用本扩展时需注意以下限制仅支持最新快照该扩展依赖 Delta Kernel API只能读取 Delta 表的最新快照无法读取任意历史快照任意快照读取能力由上游跟踪见 delta-io/delta 的 issue 2581。非分区列过滤是 best-effort对非分区列应用过滤器时可能摄入不匹配的数据详见上文过滤的语义与保证。and/or组合受限当前实现要求恰好两个谓词无法表达超过两个条件的组合除非嵌套and/or从代码结构看嵌套是可行的因为每个过滤器本身也是DeltaFilter。数据格式固定为 Parquet输入源不接收外部inputFormatDelta 表底层文件必须是版本化 Parquet 文件。版本要求需要 Delta Lake 3.0.0与 Spark 3.5.x 兼容更旧版本不受支持。深入阅读源码与测试若想进一步验证上述行为可参考仓库中的以下位置输入源实现DeltaInputSource.java、DeltaInputSourceReader.java、DeltaInputRow.java、DeltaSplit.java过滤器实现DeltaFilter.java 及filter包下的DeltaAndFilter、DeltaOrFilter、DeltaNotFilter、DeltaEqualsFilter、DeltaGreaterThanFilter、DeltaGreaterThanOrEqualsFilter、DeltaLessThanFilter、DeltaLessThanOrEqualsFilter、DeltaFilterUtils扩展装配DeltaLakeDruidModule.java测试用例extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/下的DeltaInputSourceTest、DeltaInputSourceSerdeTest、RowSerdeTest、DeltaTimeUtilsTest以及各过滤器测试其中PartitionedDeltaTable与NonPartitionedDeltaTable分别构造了分区/非分区表场景用于验证过滤行为输入源文档input-sources.md 的 Delta Lake input source 小节扩展文档delta-lake.md综上druid-deltalake-extensions是一个基于 Delta Kernel 的轻量连接器它让 Druid 得以以标准delta输入源消费 Delta Lake 表的最新快照通过扫描级过滤和列剪枝控制摄入规模并通过可拆分输入源支撑并行批式摄入。理解其仅最新快照分区列过滤才保证生效等边界是把它稳定用于生产摄入任务的关键。赞分享数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载相关推荐Apache Druid Thrift 扩展实战从实时流到 Hadoop 批量的 Thrift 数据摄取与解析Apache Druid Thrift 扩展实战从实时流到 Hadoop 批量的 Thrift 数据摄取与解析 Apache Druid 的 druid th数据库数据分析OLAP大数据实时分析数据仓库后端Apache Druid PostgreSQL 元数据存储与 PostgreSQL 批量摄入实战指南Apache Druid PostgreSQL 元数据存储与 PostgreSQL 批量摄入实战指南 Apache Druid 的 Coordinator、Ov数据库OLAP大数据后端用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务 Lucky 是一款面向软硬路由的公网管理工具集成端口转发、动态域名DDNS、反向代理、网络后端网络通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
