首页/新闻资讯/正文详情

Apache Druid 教程:使用 transformSpec 在摄取阶段转换与过滤输入数据

发布时间:2026/9/23 20:49:51 来源:云帆数科 栏目:资讯中心
Apache Druid 教程:使用 transformSpec 在摄取阶段转换与过滤输入数据
数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载本教程演示如何利用 Apache Druid 摄取规范ingestion spec中的transformSpec在数据摄入阶段对输入数据进行行级过滤与字段转换从而在数据落地为 Segment 之前完成清洗与预处理。读完本文后你将掌握transformSpec中表达式转换expression transform与过滤filter的完整配置方法、二者的执行顺序与作用机制并能够基于 Druid 表达式语言 编写出可复用的摄取期 ETL 逻辑。本教程假定你已经按照 单机快速入门 下载并启动了 Apache Druid 的本地环境并建议先完成 加载文件 与 查询数据 两个教程以便熟悉任务提交与查询的基本操作。背景transformSpec 在摄取流水线中的位置在 Druid 中一条输入记录被读取后会按照固定的顺序依次经过摄取规范中的各个组件。根据 摄取规范文档 的说明其处理顺序为flattenSpec如有用于展平嵌套数据timestampSpec解析时间戳transformSpec转换与过滤即本文主题dimensionsSpec与metricsSpec维度与指标的定义transformSpec位于dataSchema之下是可选配置由transforms转换列表与filter过滤条件两部分组成。它负责在摄入期完成两类工作转换基于输入行的字段计算新的字段值可以原地改写已有字段也可以生成全新字段过滤按条件丢弃不满足要求的输入行只有通过过滤的行才会被写入 Segment。示例数据仓库的 examples/quickstart/tutorial/transform-data.json 提供了本教程的样例数据内容如下每行一条 JSON 记录{timestamp:2018-01-01T07:01:35Z,animal:octopus, location:1, number:100} {timestamp:2018-01-01T05:01:35Z,animal:mongoose, location:2,number:200} {timestamp:2018-01-01T06:01:35Z,animal:snake, location:3, number:300} {timestamp:2018-01-01T01:01:35Z,animal:lion, location:4, number:300}每条记录包含四个字段timestampISO 格式时间戳、animal动物名称字符串、location数值、number数值。我们将通过 transformSpec 对这些数据做以下处理给animal列的值统一加上super-前缀用number * 3生成一个新的triple-number列过滤掉不满足条件的最后一行lion。完整的摄取规范transform spec 实战下面是本教程使用的完整摄取规范。你可以直接在 Druid 分发包的quickstart/tutorial/目录下找到仓库对应的实际文件 examples/quickstart/tutorial/transform-index.json下文说明其与文档展示版本的一处细微差异{ type : index_parallel, spec : { dataSchema : { dataSource : transform-tutorial, timestampSpec: { column: timestamp, format: iso }, dimensionsSpec : { dimensions : [ animal, { name: location, type: long } ] }, metricsSpec : [ { type : count, name : count }, { type : longSum, name : number, fieldName : number }, { type : longSum, name : triple-number, fieldName : triple-number } ], granularitySpec : { type : uniform, segmentGranularity : week, queryGranularity : minute, intervals : [2018-01-01/2018-01-03], rollup : true }, transformSpec: { transforms: [ { type: expression, name: animal, expression: concat(super-, animal) }, { type: expression, name: triple-number, expression: number * 3 } ], filter: { type:or, fields: [ { type: selector, dimension: animal, value: super-mongoose }, { type: selector, dimension: triple-number, value: 300 }, { type: selector, dimension: location, value: 3 } ] } } }, ioConfig : { type : index_parallel, inputSource : { type : local, baseDir : quickstart/tutorial, filter : transform-data.json }, inputFormat : { type :json }, appendToExisting : false }, tuningConfig : { type : index_parallel, partitionsSpec: { type: dynamic }, maxRowsInMemory : 25000 } } }该规范创建名为transform-tutorial的数据源使用index_parallel批式并行摄取从本地目录读取 JSON 文件。注意metricsSpec中既保留了原始number列longSum聚合也对转换生成的triple-number列做了同样的求和聚合——这正是同时摄入原始列与转换列的典型用法。版本差异说明文档展示的规范中granularitySpec.segmentGranularity为week、tuningConfig使用partitionsSpec.type: dynamic而仓库中实际落盘的示例文件 transform-index.json 将segmentGranularity设为day、并在tuningConfig中使用maxRowsPerSegment: 5000000。两者对数据转换与过滤的核心逻辑完全一致仅 Segment 划分粒度与分区调优参数不同均可直接提交运行。两个表达式转换transformstransforms列表中的每个条目都是一个表达式转换语法为{ type: expression, name: 输出字段名, expression: Druid 表达式 }本示例定义了两个转换animal原地改写字段。concat(super-, animal)会在animal列的每个值前拼接super-前缀。由于转换的name与输入字段同名都是animal转换结果会覆盖shadow原字段等价于就地变换。triple-number生成新字段。number * 3将number列的值乘以 3输出到新的triple-number列中。原始number列依然保留因此最终同时存在原始值与变换值两个字段。根据 摄取规范文档 的定义transforms列表中的字段可以被dimensionsSpec、metricsSpec等后续组件引用。需要特别注意的是转换存在两个限制转换只能引用输入行中实际存在的字段不能引用其他转换的输出多个转换之间彼此独立、顺序无关转换只能新增字段不能删除字段不过你可以用用全 null 值覆盖某个字段的方式来达到近似删除的效果。过滤条件filterfilter使用 Druid 标准的查询过滤器语法可参考 查询过滤器文档本例是一个由三个selector条件组成的or逻辑或过滤器animal等于super-mongoose注意这里匹配的是转换后的值因为过滤器在转换之后执行triple-number等于300匹配转换产物字段location等于3。三个条件中任一满足即保留该行。逐行核对样例数据原始行转换后 animaltriple-numberlocation命中的条件是否保留octopussuper-octopus3001triple-number300保留mongoosesuper-mongoose6002animalsuper-mongoose保留snakesuper-snake9003location3保留lionsuper-lion9004无丢弃最终恰好选中前 3 行最后一行lion被过滤掉。这里有一个关键语义filter 在 transforms 之后应用因此过滤条件既可以引用原始字段也可以引用转换产生的字段——这正是本示例中animalsuper-mongoose与triple-number300能够生效的原因。提交摄取任务使用分发包自带的批式任务提交脚本bin/post-index-task将任务 POST 到 Overlord默认端口 8081该脚本会持续轮询直到数据可查询bin/post-index-task --file quickstart/tutorial/transform-index.json --url http://localhost:8081如果你将上面的规范保存为自己的文件把--file指向对应路径即可。任务成功完成后数据即已写入transform-tutorial数据源。查询转换后的数据启动分发包中的 SQL 客户端bin/dsql执行查询dsql select * from transform-tutorial;预期结果如下┌──────────────────────────┬────────────────┬───────┬──────────┬────────┬───────────────┐ │ __time │ animal │ count │ location │ number │ triple-number │ ├──────────────────────────┼────────────────┼───────┼──────────┼────────┼───────────────┤ │ 2018-01-01T05:01:00.000Z │ super-mongoose │ 1 │ 2 │ 200 │ 600 │ │ 2018-01-01T06:01:00.000Z │ super-snake │ 1 │ 3 │ 300 │ 900 │ │ 2018-01-01T07:01:00.000Z │ super-octopus │ 1 │ 1 │ 100 │ 300 │ └──────────────────────────┴────────────────┴───────┴──────────┴────────┴───────────────┘ Retrieved 3 rows in 0.03s.从结果中可以验证本教程的全部要点lion行已被丢弃原始 4 行数据只剩 3 行animal列已被转换所有值都带上了super-前缀原始列与转换列并存number保留原始值如 100triple-number为转换值如 300时间戳按分钟粒度截断queryGranularity: minute将07:01:35Z归整为07:01:00Z且由于rollup: true同分钟内相同维度/指标组合的行会被合并本例中每行时间互不相同因此count均为 1。源码级原理transformSpec 是如何工作的理解了使用方式之后深入 Druid 源码可以更清晰地把握其执行语义。所有相关实现都位于 processing 模块的 org.apache.druid.segment.transform 包 下。TransformSpec过滤器 转换列表的容器TransformSpec.java 是transformSpec的配置模型持有filterDimFilter与transformsTransform 列表两个部分。其构造函数在解析规范时会校验转换名称不能重复一旦发现两个转换使用相同name会直接抛出ISE(Transform name %s cannot be used twice, ...)从源头杜绝了字段覆盖的歧义。Transform是一个标注了ExtensionPoint的扩展点接口见 Transform.java通过JsonSubTypes注册了唯一的type: expression实现——即 ExpressionTransform。ExpressionTransform接收name与expression两个参数并在构造时通过ExprMacroTable表达式宏表负责注册 Druid 内置表达式函数将表达式文本解析为可执行的Expr对象解析过程使用Suppliers.memoize惰性缓存避免重复解析。Transformer先转换、后过滤的执行引擎真正执行转换与过滤的核心是 Transformer.java 的transform(InputRow)方法其流程清晰印证了文档语义若行内配置了转换将原始行包装为TransformedInputRow持有名称到RowFunction的映射将转换后的行放入ThreadLocal交给ValueMatcher由filter编译而来进行匹配若匹配失败则返回null调用方据此丢弃该行。也就是说转换先行、过滤在后过滤天然可以引用转换产物。Transformer还提供了针对InputRowListPlusRawValues的重载支持在采样与并行摄取index_parallel场景下保留原始 raw value 的同时完成批量转换与过滤。在整条摄取链路上TransformSpec.decorate(...)提供两种接入方式见 TransformSpec.java对旧的InputRowParser体系包装为TransformingInputRowParser对新的InputSourceReader体系包装为TransformingInputSourceReader后者直接将Transformer应用到读取到的每一行上从而兼容不同代际的批式与流式摄取框架。TransformedInputRowshadow 语义与 __time 转换的落点TransformedInputRow.java 实现了 shadow 语义getDimension、getRaw、getMetric在取值时都优先检查该列是否存在对应的转换函数存在则计算转换值否则回退到原始行——这正是同名转换覆盖原字段、且转换表达式内部仍能引用被覆盖的原始字段的实现基础。此外它还支持对时间列做转换readTimestampFromRow第 56-71 行检测转换列表中是否存在名为__time的转换若存在则用其计算结果作为行时间戳。这在需要基于其他字段推导时间、或对时间做偏移的场景非常有用。Transformer构造时还会通过TransformSpec.getRequiredColumns()第 114-127 行汇总 filter 与所有转换引用的输入列供上层做列裁剪与依赖分析。测试用例佐证processing 模块的单元测试对上述行为提供了直接验证TransformSpecTest.testTransformOverwriteField验证转换允许覆盖字段、且表达式可以引用被覆盖的字段本身concat(x, y)覆盖x得到foobarTransformSpecTest.testFilterOnTransforms验证过滤器可以引用转换字段——AndDimFilter中同时使用原始字段x与转换字段f、g其中一行因不满足条件被过滤为nullTransformSpecTest.testTransformTimeFromOtherFields验证用(a b) * 3600000这类表达式从普通字段推导__time的用法TransformerTest.testTransformTimeColumn验证timestamp_shift(__time, P1D, -2)这类时间偏移转换会真实改变行的时间戳。这些测试从代码层面确认了本文描述的所有语义可以作为深入学习时的参考。扩展应用更丰富的表达式转换transforms中的expression使用完整的 Druid 表达式语言。该语言支持常规运算符*、/、%、、-、比较与逻辑运算符等具体优先级见文档并提供大量内置函数常见的摄取期用法包括字符串处理upper(country)、lower(...)、concat(...)、replace(...)、substring(...)、strlen(...)、regexp_extract(...)等可用于统一大小写、拼接、抽取子串数值计算number * 3、abs(...)、ceil(...)、floor(...)、pow(...)等可用于单位换算、取整、派生指标时间处理timestamp_parse(...)、timestamp_format(...)、timestamp_floor(...)、timestamp_shift(...)等可在转换阶段重新解析或调整时间逻辑控制if(predicate, then, else)、case_searched(...)、coalesce(...)、nvl(...)、isnull(...)等可实现条件赋值与空值兜底查询期查找表lookup(expr, lookup-name, [replaceMissingValueWith])可在摄入时直接引用已注册的 Lookup 完成字典映射。结合本文的concat(super-, animal)与number * 3两个例子你已经掌握了在 Druid 摄取阶段完成字段清洗、派生与过滤的完整套路。将这些转换逻辑前置到摄入期可以显著降低查询期的计算负担让 Segment 中存储的就是干净、可直接消费的数据。总结transformSpec位于dataSchema下由transforms与filter组成在timestampSpec之后、dimensionsSpec/metricsSpec之前执行表达式转换通过type: expressionnameexpression定义同名转换覆盖原字段shadow新名转换生成新列且转换不能引用其他转换、不能删除字段filter使用标准查询过滤器语法在转换完成后执行因此可以直接引用转换产物批式摄取可用bin/post-index-task提交任务用bin/dsql查询验证转换与过滤结果底层由 Transformer 实现先转换后过滤TransformedInputRow 实现 shadow 语义并支持__time转换相关语义均有单元测试覆盖。赞分享数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载相关推荐Apache Druid 数组类型ARRAY完全指南摄入、过滤与分组实战Apache Druid 数组类型ARRAY完全指南摄入、过滤与分组实战 Apache Druid 支持 SQL 标准的 ARRAY 类型列涵盖 VAR数据库OLAP大数据后端Apache Druid实时数据摄入与处理机制Apache Druid实时数据摄入与处理机制 Apache Druid采用独特的消防部门架构实现高效实时数据摄入通过FireDepartment、Fir数据库数据分析OLAP大数据实时分析数据仓库后端Apache Druid 多阶段查询MSQ任务引擎SQL 批量摄取完整指南Apache Druid 多阶段查询MSQ任务引擎SQL 批量摄取完整指南 本文围绕 Apache Druid 内置的 druid multi stage数据库OLAP大数据后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

用 AAS 的 cc-skill-project-guidelines-example 模板,为真实项目编写项目专属 Skill
用 AAS 的 cc-skill-project-guidelines-example 模板,为真实项目编写项目专属 Skill

AI 技能AI 插件 【免费下载链接】agentic-awesome-skills AAS Core is the local, agent-first control plane for complete catalog discovery, agent-owned selection, stack validation, and planning, backed by 2,445 agentic skills. Includes CLI, local MCP, catalog, … · 2026/9/23 20:49:44

Dopamine 实验数据工具集:dopamine.colab.utils 源码级解析与实战
Dopamine 实验数据工具集:dopamine.colab.utils 源码级解析与实战

Dopamine 实验数据工具集:dopamine.colab.utils 源码级解析与实战 【免费下载链接】dopamine Dopamine is a research framework for fast prototyping of reinforcement learning algorithms. 项目地址: https://gitcode.com/gh_mirrors/do/dopamine dopam… · 2026/9/23 20:49:44

asfd面试必问:3分钟搞定市政公用工程与游戏开发选型
asfd面试必问:3分钟搞定市政公用工程与游戏开发选型

asfd面试必问:3分钟搞定市政公用工程与游戏开发选型 翻开官方文档想搞懂 asfd,结果目录比书还厚,翻到第三页就懵了?别慌,这正是很多老手都会遇到的死胡同。其实 asfd… · 2026/9/23 20:49:44

GB0-670 H3C MSA存储认证备考:架构、配置与实战故障处理
GB0-670 H3C MSA存储认证备考:架构、配置与实战故障处理

简介:H3CNE-MSA(GB0-670)认证题库文档,面向备考H3C认证工程师-MSA存储的考生和技术服务人员,覆盖MSA 2042/2040设备差异、主机访问协议、存储扩展、备份技术、SMU管理、SAN/NAS、虚拟存储、多路径、SSD优势等核心考点。… · 2026/9/23 21:29:57

spotifyd 开发环境搭建与代码贡献完整指南:从编译运行到提交 PR
spotifyd 开发环境搭建与代码贡献完整指南:从编译运行到提交 PR

音频后端 【免费下载链接】spotifyd A spotify daemon 项目地址: https://gitcode.com/gh_mirrors/sp/spotifyd 点击查看 免费下载 导读 本文基于 spotifyd 仓库根目录的 CONTRIBUTING.md 展开,面向希望为 spotifyd 贡献代码或亲自从源码编译运行的开发… · 2026/9/23 21:29:51

医疗知识图谱构建与KBQA问答系统实战:从实体识别到Neo4j查询
医疗知识图谱构建与KBQA问答系统实战:从实体识别到Neo4j查询

简介:这是一套面向医疗领域知识图谱问答(KBQA)系统从零构建的完整资料包,适合希望快速上手知识图谱与智能问答的AI开发者、算法工程师及高校学生。项目包含7类实体、约3.7万实体、21万实体关系的医疗知识图谱构建案例,… · 2026/9/23 21:29:44

GB0-670备考指南:从MSA存储架构到双控切换实战解析
GB0-670备考指南:从MSA存储架构到双控切换实战解析

简介:面向H3CNE-MSA认证备考者的Word版题库,聚焦H3C代理的MSA存储设备基础配置与维护技术。文档以单个docx文件封装,大小约32KB,包含大量单选与多选试题,内容覆盖MSA 2040/2042产品特性、iSCSI与SAS等主机访问协议、磁… · 2026/9/23 21:29:44

QUANTAXIS 数据流处理与事件驱动架构深度解析:从迭代器到分布式消息队列
QUANTAXIS 数据流处理与事件驱动架构深度解析:从迭代器到分布式消息队列

金融科技后端数据分析 【免费下载链接】QUANTAXIS QUANTAXIS 支持任务调度 分布式部署的 股票/期货/期权 数据/回测/模拟/交易/可视化/多账户 纯本地量化解决方案 项目地址: https://gitcode.com/gh_mirrors/qu/QUANTAXIS 点击查看 免费下载 导读:本文聚… · 2026/9/23 21:29:44

服务器运行报告模板自动化生成与监控指标设计指南
服务器运行报告模板自动化生成与监控指标设计指南

简介:服务器运行报告模板是一份专为IT运维人员设计的标准化文档,用于日常服务器巡检、定期维护与故障排查记录。模板涵盖设备硬件信息、机柜防尘与风扇噪音检查、电源与硬盘状态、操作系统及应用程序运行状况,并给出内存、CPU、硬盘、系统信息… · 2026/9/23 21:29:38

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码