
数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载本文介绍如何在 Apache Druid Web Console 的Query视图中将已有的原生批量摄入native batch ingestionspec 一键转换为druid-multi-stage-queryMSQ扩展支持的 SQL 查询任务。读完本文你将掌握转换入口的位置、支持与不支持的 spec 形态、转换器逐项规则的底层实现时间戳、维度、指标、分区、写入策略等并能熟练复核、微调并执行生成的 SQL 完成数据摄入。背景为什么需要把 spec 转成 SQLApache Druid 24.0 起内置了druid-multi-stage-query扩展它为 SQL 增加了多阶段查询MSQ任务引擎可以把INSERT、REPLACE等 SQL 语句当作批量任务在索引服务上执行并像其他批量摄入方式一样发布 segment。相比传统 JSON specSQL-based ingestion 提供了更灵活的表达能力例如在INSERT ... SELECT/REPLACE ... SELECT中直接施加转换、过滤、JOIN 与聚合甚至支持基于查询结果建新表in-database transformation。如果你已经在使用原生批量摄入index_parallel等任务类型不必从零手写 SQL——Web Console 提供了Convert ingestion spec to SQL功能可以把现成的 ingestion spec 自动翻译为等价的 SQL 查询。这正是本教程要演示的核心能力。转换入口Query 视图中的菜单项该功能位于 Web Console 的Query视图。在 workbench-view.tsx 中可以看到只有当当前查询引擎支持sql-msq-task时顶部的运行菜单runMoreMenu才会渲染Convert ingestion spec to SQL菜单项点击后触发openSpecDialog打开转换对话框见 workbench-view.tsx 与 renderSpecDialog。操作步骤如下打开 Web Console 的Query视图找到包含Run按钮的菜单栏。点击菜单栏中的省略号ellipsis图标选择Convert ingestion spec to SQL。在弹出的Ingestion spec to convert窗口中粘贴你的 ingestion specJSON。你可以使用自己的 spec也可以使用下文教程附带的示例 spec。点击Submit提交 specWeb Console 会基于该 JSON spec 生成一段等价 SQL并自动打开一个新查询 TabTab 名为Convert dataSource见 workbench-view.tsx。仔细复核生成的 SQL确认它符合你的预期、与原始 spec 语义一致。点击Run启动摄入任务。如果 JSON 解析失败对话框会通过 spec-dialog.tsx 显示带行列号at line X, column Y的错误提示如果转换器本身抛错则会以 toast 形式提示Could not convert spec: ...见 workbench-view.tsx。支持哪些 spec 类型转换入口并非对任意 spec 都生效。核心转换函数convertSpecToSql位于 spec-conversion.ts它的第一道检查就限定了输入范围spec-conversion.tsif (!oneOf(spec.type, index_parallel, index, index_hadoop)) { throw new Error(Only index_parallel, index, and index_hadoop specs are supported); }即目前只支持三种任务类型index_parallel并行原生批量摄入最常用index单机原生批量摄入index_hadoopHadoop 批量摄入提交后转换器会先调用upgradeSpec定义于 ingestion-spec.tsx对旧式 spec 做规范化包括把ioConfig.firehose迁移为inputSourcestatic-s3→s3、static-google-blobstore→google把旧版dataSchema.parser分解为timestampSpecdimensionsSpecinputFormat。因此较早期版本的 spec 也可以直接粘贴进来。另外需要说明spatialDimensions空间维度当前无法在 SQL-based ingestion 中使用转换器会直接抛错spec-conversion.ts。完整示例从 spec 到 SQL 的转换全过程以下是教程附带的示例 ingestion spec它从https://druid.apache.org/data/wikipedia.json.gz读取 JSON 数据加载到名为wikipedia的表{ type: index_parallel, spec: { ioConfig: { type: index_parallel, inputSource: { type: http, uris: [ https://druid.apache.org/data/wikipedia.json.gz ] }, inputFormat: { type: json } }, tuningConfig: { type: index_parallel, partitionsSpec: { type: dynamic } }, dataSchema: { dataSource: wikipedia, timestampSpec: { column: timestamp, format: iso }, dimensionsSpec: { dimensions: [ isRobot, channel, flags, isUnpatrolled, page, diffUrl, { type: long, name: added }, comment, { type: long, name: commentLength }, isNew, isMinor, { type: long, name: delta }, isAnonymous, user, { type: long, name: deltaBucket }, { type: long, name: deleted }, namespace, cityName, countryName, regionIsoCode, metroCode, countryIsoCode, regionName ] }, granularitySpec: { queryGranularity: none, rollup: false, segmentGranularity: day } } } }点击Submit后Web Console 生成的 SQL 如下-- This SQL query was auto generated from an ingestion spec REPLACE INTO wikipedia OVERWRITE ALL WITH source AS (SELECT * FROM TABLE( EXTERN( {type:http,uris:[https://druid.apache.org/data/wikipedia.json.gz]}, {type:json}, [{name:timestamp,type:string},{name:isRobot,type:string},{name:channel,type:string},{name:flags,type:string},{name:isUnpatrolled,type:string},{name:page,type:string},{name:diffUrl,type:string},{name:added,type:long},{name:comment,type:string},{name:commentLength,type:long},{name:isNew,type:string},{name:isMinor,type:string},{name:delta,type:long},{name:isAnonymous,type:string},{name:user,type:string},{name:deltaBucket,type:long},{name:deleted,type:long},{name:namespace,type:string},{name:cityName,type:string},{name:countryName,type:string},{name:regionIsoCode,type:string},{name:metroCode,type:string},{name:countryIsoCode,type:string},{name:regionName,type:string}] ) )) SELECT TIME_PARSE(timestamp) AS __time, isRobot, channel, flags, isUnpatrolled, page, diffUrl, added, comment, commentLength, isNew, isMinor, delta, isAnonymous, user, deltaBucket, deleted, namespace, cityName, countryName, regionIsoCode, metroCode, countryIsoCode, regionName FROM source PARTITIONED BY DAY对照原文 spec 逐段解读REPLACE INTO wikipedia OVERWRITE ALL目标表为wikipedia采用整表覆盖写入关于 REPLACE 的语义可参考 multi-stage query reference。原 spec 没有配置dropExisting与intervals因此默认OVERWRITE ALL。EXTERN(...)第一、二参数分别序列化自ioConfig.inputSource与ioConfig.inputFormat第三个参数是根据dimensionsSpectimestampSpec生成的列类型声明数组。关于EXTERN的细节见 concepts.md。TIME_PARSE(timestamp) AS __time由timestampSpeccolumn: timestamp、format: iso推导而来。PARTITIONED BY DAY由granularitySpec.segmentGranularity: day推导而来。转换规则逐项解析源码级生成的 SQL 并非简单拼装而是由 spec-conversion.ts 中的规则逐项翻译而成。理解这些规则才能在转换后高效复核并微调结果。1. 时间戳列与__time表达式转换器要求spec.dataSchema.timestampSpec必须存在spec-conversion.ts并根据timestampSpec.format选择不同的__time表达式spec-conversion.tstimestampSpec.format生成的__time表达式列类型声明autoCASE WHEN CAST(col AS BIGINT) 0 THEN MILLIS_TO_TIMESTAMP(CAST(col AS BIGINT)) ELSE TIME_PARSE(TRIM(col)) ENDVARCHARisoTIME_PARSE(col)VARCHARposixMILLIS_TO_TIMESTAMP(col * 1000)BIGINTmillisMILLIS_TO_TIMESTAMP(col)BIGINTmicroMILLIS_TO_TIMESTAMP(col / 1000)BIGINTnanoMILLIS_TO_TIMESTAMP(col / 1000000)BIGINT其他自定义格式TIME_PARSE(col, format)VARCHAR边界情况若timestampSpec.column为__time本身则直接沿用该列spec-conversion.ts若timestampSpec.column为!!!_no_such_column_!!!即NO_SUCH_COLUMN则使用timestampSpec.missingValue构造常量时间缺省时退化为TIMESTAMP 1970-01-01spec-conversion.ts若配置了timestampSpec.missingValue会在时间表达式外包裹COALESCE(..., TIME_PARSE(missingValue))spec-conversion.ts。2. 查询粒度queryGranularity的折叠granularitySpec.queryGranularity会被应用到时间表达式上spec-conversion.ts映射表QUERY_GRANULARITY_MAP定义在 spec-conversion.tsqueryGranularity时间表达式变换none原表达式不变secondTIME_FLOOR(expr, PT1S)minuteTIME_FLOOR(expr, PT1M)fifteen_minuteTIME_FLOOR(expr, PT15M)thirty_minuteTIME_FLOOR(expr, PT30M)hourTIME_FLOOR(expr, PT1H)dayTIME_FLOOR(expr, P1D)weekTIME_FLOOR(expr, P7D)monthTIME_FLOOR(expr, P1M)quarterTIME_FLOOR(expr, P3M)yearTIME_FLOOR(expr, P1Y)无法识别的 queryGranularity 会直接抛错spec.dataSchema.granularitySpec.queryGranularity is not recognizedspec-conversion.ts。3. 维度声明与 EXTERN 列类型dimensionsSpec.dimensions中的每个维度会经inflateDimensionSpec规范化后由dimensionSpecToSqlType映射为 SQL 类型spec-conversion.ts普通字符串维度 →VARCHAR显式long/float/double维度 → 对应数值类型json类型维度 →COMPLEXjsonauto类型且带castToType→ 按目标类型转换。最终所有列声明组成EXTERN的第三个参数EXTEND (...)子句。需要留意如果metricsSpec存在其fieldName对应的列也会被追加到列声明中spec-conversion.ts并做去重dedupe。4. 指标metricsSpec的聚合转换metricSpecToSqlExpressionspec-conversion.ts给出了完整的指标映射规则指标类型生成的 SQL 表达式countCOUNT(*)longSum/floatSum/doubleSumSUM(col)longMin/floatMin/doubleMinMIN(col)longMax/floatMax/doubleMaxMAX(col)doubleFirst/floatFirst/longFirstEARLIEST(col)stringFirstEARLIEST(col, maxStringBytes ?? 128)doubleLast/floatLast/longLastLATEST(col)stringLastLATEST(col, maxStringBytes ?? 128)thetaSketchAPPROX_COUNT_DISTINCT_DS_THETA(col, size ?? 16384)HLLSketchBuild/HLLSketchMergeAPPROX_COUNT_DISTINCT_DS_HLL(col, lgK ?? 12, tgtHllType ?? HLL_4)quantilesDoublesSketchDS_QUANTILES_SKETCH(col, k ?? 128)hyperUniqueAPPROX_COUNT_DISTINCT_BUILTIN(col)已知不支持返回 undefined 并输出-- could not convert metric: ...注释的类型包括tDigestSketch、momentSketch、fixedBucketsHistogramspec-conversion.ts。指标表达式输出形如SUM(added) AS sum_addedmetricSpecToSelectspec-conversion.ts。5. rollup 与 GROUP BYgranularitySpec.rollup默认视为true?? truespec-conversion.ts。当启用 rollup 时生成的 SQL 会追加GROUP BY 各维度表达式的序号从而让指标聚合生效禁用 rollup 时则不生成 GROUP BYspec-conversion.ts。6. 段粒度与分区PARTITIONED BY由segmentGranularity生成all→ALL TIME其余值转大写如day→DAYspec-conversion.ts。若tuningConfig.partitionsSpec配置了partitionDimensions或旧的partitionDimension则追加CLUSTERED BY col1, col2, ...spec-conversion.ts对应原生摄入的single_dim分区。7. 写入目标与覆盖策略写入语句由ioConfig与数据源决定spec-conversion.ts默认生成REPLACE INTO dataSource OVERWRITE ALL若ioConfig.appendToExisting为 true则生成INSERT INTO dataSource追加语义若ioConfig.dropExisting为 true则生成REPLACE INTO dataSource OVERWRITE WHERE __time 时间区间条件区间来自granularitySpec.intervals。8. 输入源的三种形态index_parallel/index读取ioConfig.inputSource经 JSON 序列化后填入EXTERN的第一个参数spec-conversion.tsindex_hadoop读取ioConfig.inputSpec且仅支持type: static并按路径前缀映射输入源类型s3://→s3、gs://→google、hdfs://→hdfs否则抛unsupportedspec-conversion.ts输入源为druid类型从已有 datasource 读取时不生成EXTERN而是生成WITH source AS (SELECT * FROM table WHERE __time IN interval)spec-conversion.ts。9. 自动生成的 Query Context转换结果除queryString外还会携带一个queryContextspec-conversion.tsWeb Console 会把它一并写入新 Tab见 workbench-view.tsx。上下文映射规则spec 配置生成的 query contexttuningConfig.maxNumConcurrentSubTasks 1maxNumTasks maxNumConcurrentSubTasks 1控制器 工作节点tuningConfig.maxParseExceptions数值maxParseExceptionstuningConfig.indexSpecindexSpec沿用原索引压缩配置多值列以数组模式摄入getArrayMode返回arraysarrayIngestMode array前两条映射见 spec-conversion.tsindexSpec见 spec-conversion.tsarrayIngestMode见 spec-conversion.ts。在 spec-conversion.spec.ts 的测试中可以看到maxNumConcurrentSubTasks: 4与maxParseExceptions: 3、indexSpec.dimensionCompression: lzf分别被转换为maxNumTasks: 5、maxParseExceptions: 3、indexSpec: { dimensionCompression: lzf }。转换局限与人工修正提示并非 spec 中的每个元素都能机械翻译。转换器对无法自动转换的部分会在生成的 SQL 中以--:ISSUE:注释显式标出提醒你手动处理transform 转换transformSpec.transforms中的表达式无法自动翻译为 SQL 表达式。涉及维度与__time的 transform 会分别生成--:ISSUE: Transform for dimension could not be converted/--:ISSUE: Transform for __time could not be converted注释spec-conversion.ts。filter 转换仅selector、in、not、and、or五种过滤器类型可自动转换为 SQLWHERE条件convertFilterspec-conversion.ts其他类型会回退为WHERE REWRITE_[filter json]_TO_SQL --:ISSUE: The spec contained a filter that could not be automatically converted, please convert it manuallyspec-conversion.ts。不受支持的指标类型如tDigestSketch、momentSketch、fixedBucketsHistogram会输出-- could not convert metric: ...注释spec-conversion.ts。spatialDimensions直接抛错无法转换。因此转换后务必逐行复核 SQL检查--:ISSUE:注释涉及的 transform/filter 是否需手工改写确认分区、覆盖策略与 rollup 语义符合预期再点击Run。也可以先用Explain SQL query或小型样例数据验证语义菜单项在 workbench-view.tsx。测试与验证转换器是如何被保障的该转换功能有完整的单元测试覆盖位于 spec-conversion.spec.ts并通过快照toMatchSnapshot校验生成的 SQL 文本。测试场景包括index_parallelspec无 rollup含not/selector过滤器、single_dim分区、indexSpec、上下文参数映射spec-conversion.spec.tsindex_parallelspec有 rollup含count、longSum、longMax、thetaSketch等指标验证GROUP BY与APPROX_COUNT_DISTINCT_DS_THETA的生成spec-conversion.spec.tsindex_hadoopspec有 rollup验证旧式parser/inputSpec的规范化与s3://输入源映射spec-conversion.spec.ts带__timetransform / 维度 transform / 未知 filter 的 spec验证--:ISSUE:标记的输出spec-conversion.spec.tsARRAY 模式验证arrayIngestMode array上下文与castToType: ARRAY...列声明spec-conversion.spec.ts。这些测试恰好也是理解转换器行为边界的最佳参考即使你手头没有 Web Console也可以通过阅读这些用例确认某个 spec 特性会被翻译成什么 SQL。常见问题与下一步Q转换后点击 Run 报 403 错误A使用EXTERN以及S3、LOCALFILES等输入源专属表函数需要具备资源类型EXTERNAL、资源名EXTERNAL上的 READ 权限详见 multi-stage query 文档。请检查集群的授权配置。Q生成的 SQL 与我的原生 spec 不完全等价A转换是规则化的机械翻译存在上文列举的局限transform、部分 filter、部分 sketch 指标。SQL-based ingestion 的完整能力与限制可查阅 concepts.md、reference.md 与 known-issues.md更多可直接运行的示例见 examples.md。如需了解各摄入方式的适用场景可参考摄入方式对比表。Q如何判断我的集群是否启用了 MSQ 引擎A确保druid.extensions.loadlist中包含druid-multi-stage-query并在common.runtime.properties中完成配置multi-stage query 文档Web Console 的 Query 视图也只有在检测到sql-msq-task引擎时才会显示Convert ingestion spec to SQL菜单。通过本文的步骤你可以把已有的原生批量摄入 spec 快速迁移到 SQL-based ingestion在保留原有数据模型与分区策略的同时获得 SQL 的表达力与 MSQ 任务引擎的批处理能力。赞分享数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载相关推荐Apache Druid 基于 Hadoop 的批量摄入Hadoop-based Ingestion完整实战指南Apache Druid 基于 Hadoop 的批量摄入Hadoop based Ingestion完整实战指南 本文是 Apache Druid 官方参考数据库OLAP大数据后端Apache Druid 多阶段查询MSQ任务引擎SQL 批量摄取完整指南Apache Druid 多阶段查询MSQ任务引擎SQL 批量摄取完整指南 本文围绕 Apache Druid 内置的 druid multi stage数据库OLAP大数据后端Apache Druid 流式数据摄入实战基于 Tranquility Server 编写自定义 Ingestion SpecApache Druid 流式数据摄入实战基于 Tranquility Server 编写自定义 Ingestion Spec 导读 本文面向已按 singl数据库数据分析OLAP大数据实时分析数据仓库后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考