ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Druid 批式数据摄入实战:Hadoop 批式摄入任务(Batch Ingestion)完整指南

Apache Druid 批式数据摄入实战:Hadoop 批式摄入任务(Batch Ingestion)完整指南 数据库数据分析OLAP大数据实时分析数据仓库后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid7/druid点击查看免费下载Druid 的批式摄入Batch Ingestion用于将静态文件如 HDFS、S3 上的数据文件批量加载进数据仓库是构建离线数仓与历史数据回填的核心入口。本文以 batch-ingestion.md 为主体结合仓库内indexing-hadoop、indexing-service的源码实现系统讲解 Hadoop 批式摄入任务的完整结构、配置参数、分区策略、集群适配EMR / Kerberos与替代方案命令行 Hadoop Indexer、IndexTask帮助读者直接编写可运行的任务 JSON 并理解底层执行原理。一、Druid 批式摄入的三种方式概览Druid 支持从静态文件批量加载数据主要路径包括Hadoop 批式摄入Hadoop-based Batch Ingestion通过index_hadoop任务提交到运行中的 Overlord利用 Hadoop MapReduce 集群完成大规模数据的分区、排序与 segment 构建适合海量数据、规模化摄入命令行 Hadoop IndexerCommand Line Hadoop Indexer不依赖完整的索引服务独立进程内直接运行 Hadoop 摄入详见 command-line-hadoop-indexer.mdIndexTask 批式摄入IndexTask-based Batch Ingestion不需要 Hadoop 依赖直接在 Druid 进程内完成摄入比 Hadoop 方式更慢、扩展性更低适合中小数据量详见 tasks.md。其中 Hadoop 批式摄入是本文的绝对重点任务被 POST 到 Druid Overlordindexing-service.md由 Overlord 协调调度后交给 Peon 进程执行。二、Hadoop 批式摄入任务结构剖析与源码印证index_hadoop任务是 Hadoop 摄入的载体。从源码 HadoopIndexTask.java 可以看到任务类型getType()固定返回index_hadoop其运行流程分为两个阶段确定配置阶段HadoopDetermineConfigInnerProcessing.runTask使用HadoopDruidDetermineConfigurationJob推算数据 interval、分区信息并注入 Overlord 提供的segmentOutputPath与workingPath生成索引阶段HadoopIndexGeneratorInnerProcessing.runTask通过HadoopDruidIndexerJob真正执行 MapReduce 摄入生成 segment 列表最终由 Overlord 统一发布到元数据存储。该任务还通过invokeForeignLoader加载独立的 Hadoop 依赖 ClassLoader这正是hadoopDependencyCoordinates存在的意义——用独立的 Hadoop 版本运行摄入避免与 Druid 自身依赖冲突。一个完整的任务示例摘自原文档可直接复用{ type : index_hadoop, spec : { dataSchema : { dataSource : wikipedia, parser : { type : hadoopyString, parseSpec : { format : json, timestampSpec : { column : timestamp, format : auto }, dimensionsSpec : { dimensions: [page,language,user,unpatrolled,newPage,robot,anonymous,namespace,continent,country,region,city], dimensionExclusions : [], spatialDimensions : [] } } }, metricsSpec : [ { type : count, name : count }, { type : doubleSum, name : added, fieldName : added }, { type : doubleSum, name : deleted, fieldName : deleted }, { type : doubleSum, name : delta, fieldName : delta } ], granularitySpec : { type : uniform, segmentGranularity : DAY, queryGranularity : NONE, intervals : [ 2013-08-31/2013-09-01 ] } }, ioConfig : { type : hadoop, inputSpec : { type : static, paths : /MyDirectory/example/wikipedia_data.json } }, tuningConfig : { type: hadoop } }, hadoopDependencyCoordinates: my_hadoop_version }任务顶层字段说明propertydescriptionrequired?type任务类型固定为index_hadoop。yesspecHadoop Index Spec结构见下文的 DataSchema / IOConfig / TuningConfig。yeshadoopDependencyCoordinatesJSON 数组指定 Druid 使用的 Hadoop 依赖坐标覆盖默认 Hadoop 坐标一旦指定Druid 会从druid.extensions.hadoopDependenciesDir指定位置加载这些依赖。noclasspathPrefix预附加到 Peon 进程的 classpath。no此外Druid 会自动为运行在 Hadoop 集群中的 job container 计算 classpath若 Hadoop 与 Druid 依赖冲突可手动通过druid.extensions.hadoopContainerDruidClasspath属性指定详见 configuration/index.md 的 extensions 配置。2.1 关键限制segmentOutputPath 与 workingPath 必须缺省提交到 Overlord 的index_hadoop任务与命令行 Hadoop Indexer 的 spec 有本质区别从 HadoopIndexTask.java 构造器的Preconditions.checkArgument可以确认通过索引服务运行时segmentOutputPath、workingPath、metadataUpdateSpec三个字段必须为 null/缺省——因为 segment 输出路径、工作目录和元数据更新均由 Overlord 统一管理任务运行时由HadoopDetermineConfigInnerProcessing动态注入。2.2 DataSchema该字段必填。负责定义数据源名称、解析器parser/parseSpec、指标聚合metricsSpec与粒度granularitySpec通用说明见 ingestion/index.md。任务示例中的hadoopyStringparser 对应源码 HadoopyStringInputRowParser.java它把 Hadoop 输入的字节串按 parseSpec 解析成行。三、IOConfig输入来源与 segment 输出ioConfig字段必填类型固定为hadoop对应源码类 HadoopIOConfig.java。FieldTypeDescriptionRequiredtypeString固定为hadoop。yesinputSpecObject指定数据拉取来源详见下文。yessegmentOutputPathStringsegment 输出目录路径。yes仅命令行模式索引服务模式下必须缺省metadataUpdateSpecObject指定如何为这些 segment 更新 Druid 集群的元数据。yes仅命令行模式索引服务模式下必须缺省3.1 InputSpec 的四种类型inputSpec由 Druid 通过 Jackson 多态反序列化见 PathSpec.java 的子类体系目前支持static、granularity、dataSource、multi四种static提供数据文件的静态路径对应 StaticPathSpec.java。FieldTypeDescriptionRequiredpathsArray of String指示原始数据所在位置的输入路径字符串。yes例如支持逗号分隔多路径paths : s3n://billy-bucket/the/data/is/here/data.gz, s3n://billy-bucket/the/data/is/here/moredata.gz, s3n://billy-bucket/the/data/is/here/evenmoredata.gz实现细节上StaticPathSpec.addInputPaths通过HadoopGlobPathSplitter拆分 Hadoop glob 表达式规避 HADOOP MAPREDUCE-5061 的 bug再经MultipleInputs注册输入。若未显式指定inputFormat默认使用TextInputFormat当 TuningConfig 的combineTexttrue时自动切换为CombineTextInputFormat以合并小文件。granularity按时间组织目录结构的数据目录路径格式为yXXXX/mXX/dXX/HXX/MXX/SXX日期小写、时间大写对应 GranularityPathSpec.java。FieldTypeDescriptionRequireddataGranularityString期望的数据粒度例如 hour 表示期望目录yXXXX/mXX/dXX/HXX。yesinputPathString追加时间路径的基础路径。yesfilePatternString文件需匹配才被纳入摄入的正则模式。yespathFormatString每个目录的 Joda datetime 格式默认yyyyy/mMM/ddd/HHH参见 Joda DateTimeFormat 文档。no例如若以 2012-06-01/2012-06-02 的 interval 运行将期望数据位于s3n://billy-bucket/the/data/is/here/y2012/m06/d01/H00 s3n://billy-bucket/the/data/is/here/y2012/m06/d01/H01 ... s3n://billy-bucket/the/data/is/here/y2012/m06/d01/H23从 GranularityPathSpec.java 的addInputPaths实现可以看到完整发现逻辑先用dataGranularity.getIterable(inputInterval)把摄入 interval 按粒度切成多个时间桶再对每个桶用FSSpideringIterator.spiderIterable递归遍历文件系统最后用filePattern编译的正则逐个匹配文件路径命中者才加入输入列表——这意味着文件名的过滤是纯正则匹配目录遍历是递归式的。dataSource读取已有的 Druid segment 作为输入用于增量更新详见 update-existing-data.md。其实现类为 DatasourcePathSpec.java配合 DatasourceInputFormat.java 将历史 segment 作为 MapReduce 输入。multi同时读取多个来源的数据混合 static / granularity / dataSource详见 update-existing-data.md对应 MultiplePathSpec.java。四、TuningConfig摄入调优参数全解tuningConfig可选缺省时使用默认参数。其 JSON 反序列化对应 HadoopTuningConfig.java类型固定为hadoop。FieldTypeDescriptionRequiredworkingPathStringHadoop job 之间中间结果的工作路径。no默认/tmp/druid-indexing索引服务模式下由 Overlord 注入versionString创建 segment 的版本除非useExplicitVersiontrue否则对 HadoopIndexTask 忽略。no默认摄入开始时刻partitionsSpecObject每个时间桶如何划分为 segment 的说明缺省表示不分区。no默认hashedmaxRowsInMemoryInteger聚合后落盘前的行数上限注意是 roll-up 之后的 post-aggregation 行数可能不等于输入事件数用于控制所需 JVM 堆大小。no默认 75000leaveIntermediateBooleanjob 完成无论成功失败后在工作路径保留中间文件便于调试。no默认 falsecleanupOnFailureBooleanjob 失败时清理中间文件除非开启 leaveIntermediate。no默认 trueoverwriteFilesBoolean摄入时覆盖已有文件。no默认 falseignoreInvalidRowsBoolean忽略存在问题的行。no默认 falsecombineTextBoolean使用 CombineTextInputFormat 将多个文件合并为一个 split处理大量小文件时可加速。no默认 falseuseCombinerBoolean尽可能在 mapper 端使用 Hadoop combiner 合并行。no默认 falsejobPropertiesObject追加到 Hadoop job 配置的属性映射见下文。no默认 nullindexSpecObject调整数据索引方式见下文。nobuildV9DirectlyBoolean直接构建 v9 索引而不是先构建 v8 再转换。no默认 truenumBackgroundPersistThreadsInteger用于增量 persist 的后台线程数使用会显著增加内存与 CPU 压力但加速任务建议从默认 0当前线程 persist改为 1。no默认 0forceExtendableShardSpecsBoolean强制使用 extendable shardSpecs实验性功能配合 kafka-ingestion.md 使用。no默认 falseuseExplicitVersionBoolean强制 HadoopIndexTask 使用version字段。no默认 false源码 HadoopTuningConfig.java 佐证了上述默认值DEFAULT_ROW_FLUSH_BOUNDARY 75000、默认partitionsSpec为HashedPartitionsSpec.makeDefaultHashedPartitionsSpec()、DEFAULT_BUILD_V9_DIRECTLY Boolean.TRUE、DEFAULT_NUM_BACKGROUND_PERSIST_THREADS 0且构造器断言numBackgroundPersistThreads 0。同时该源码保留了对历史字段rowFlushBoundary的兼容解析maxRowsInMemory与rowFlushBoundary二选一。4.1 jobProperties 字段tuningConfig : { type: hadoop, jobProperties: { hadoop-property-a: value-a, hadoop-property-b: value-b } }Hadoop 的 MapReduce 文档列出了所有可配置参数。部分 Hadoop 发行版可能需要设置mapreduce.job.classpath或mapreduce.job.user.classpath.first以避免类加载问题详见 other-hadoop.md。4.2 IndexSpecFieldTypeDescriptionRequiredbitmapObjectbitmap 索引的压缩格式JSON 对象选项见下文。no默认 ConcisedimensionCompressionString维度列压缩格式LZ4、LZF或uncompressed。no默认LZ4metricCompressionString指标列压缩格式LZ4、LZF、uncompressed或none。no默认LZ4longEncodingStringlong 类型指标列与维度列的编码auto按列基数用 offset 或查找表并以变长存储longs每值固定 8 字节原样存储。no默认longsBitmap 类型Concise bitmapFieldTypeDescriptionRequiredtypeString必须为concise。yesRoaring bitmapFieldTypeDescriptionRequiredtypeString必须为roaring。yescompressRunOnSerializationBoolean估计 RLE 更省空间时使用游程编码。no默认true五、分区规范Partitioning specificationhashed 与 dimensionSegment 总是先按时间戳分区依据 granularitySpec再根据分区类型进一步划分。Druid 支持两种分区策略hashed基于每行所有维度哈希与dimension基于单一维度取值区间。多数场景推荐hashed分区相比单维度分区哈希分区索引性能更好、segment 大小更均匀。5.1 哈希分区Hash-based partitioningpartitionsSpec: { type: hashed, targetPartitionSize: 5000000 }哈希分区先自动选择 segment 数量再依据每行所有维度哈希将行分布到这些 segment。segment 数量根据输入集基数与目标分区大小自动决定对应 HashedPartitionsSpec.java 与 DetermineHashedPartitionsJob.java。FieldDescriptionRequiredtype分区类型。hashedtargetPartitionSize每个分区的目标行数建议使 segment 达到 500MB~1GB。targetPartitionSize 与 numShards 二选一numShards直接指定分区数量而非目标大小跳过自动选分区数步骤摄入更快。numShards 与 targetPartitionSize 二选一partitionDimensions参与分区的维度留空表示所有维度仅与 numShards 搭配使用设置 targetPartitionSize 时被忽略。no5.2 单维度分区Single-dimension partitioningpartitionsSpec: { type: dimension, targetPartitionSize: 5000000 }单维度分区先选择分区的维度再将该维度划分为连续区间每个 segment 包含该维度取值落在对应区间内的所有行。例如可按维度 host 用区间 a.example.com~f.example.com 与 f.example.com~z.example.com 分区。默认自动选择维度也可显式指定对应 SingleDimensionPartitionsSpec.java 与 DeterminePartitionsJob.java。FieldDescriptionRequiredtype分区类型。dimensiontargetPartitionSize每个分区的目标行数建议使 segment 达到 500MB~1GB。yesmaxPartitionSize每个分区的最大行数默认比 targetPartitionSize 大 50%。nopartitionDimension用于分区的维度留空自动选择。noassumeGrouped假设输入数据已按时间与维度分组开启后摄入更快但假设被违反时可能选出次优分区。no六、运行环境适配远程集群、EMR 与 Kerberos 安全集群6.1 远程 Hadoop 集群若使用远程 Hadoop 集群务必把存放配置*.xml文件的目录放入 Druid 的_common配置目录参考 examples/conf/druid/_common/。若 Hadoop 版本与 Druid 编译版本存在依赖问题参见 other-hadoop.md。6.2 使用 Elastic MapReduceEMR从 S3 摄入集群运行在 AWS 时可用 EMR 索引 S3 数据步骤如下创建持久化、长期运行的 EMR 集群创建集群时向导模式下位于高级设置的 Edit software settings填入以下配置classificationyarn-site,properties[mapreduce.reduce.memory.mb6144,mapreduce.reduce.java.opts-server -Xms2g -Xmx2g -Duser.timezoneUTC -Dfile.encodingUTF-8 -XX:PrintGCDetails -XX:PrintGCTimeStamps,mapreduce.map.java.opts758,mapreduce.map.java.opts-server -Xms512m -Xmx512m -Duser.timezoneUTC -Dfile.encodingUTF-8 -XX:PrintGCDetails -XX:PrintGCTimeStamps,mapreduce.task.timeout1800000]按照 tutorials/cluster.md 中 Configure Hadoop for data loads 的指引使用 EMR master 上/etc/hadoop/conf中的 XML 文件完成配置。通过 EMR 从 S3 加载数据在 Hadoop 摄入任务tuningConfig的jobProperties中追加jobProperties : { fs.s3.awsAccessKeyId : YOUR_ACCESS_KEY, fs.s3.awsSecretAccessKey : YOUR_SECRET_KEY, fs.s3.impl : org.apache.hadoop.fs.s3native.NativeS3FileSystem, fs.s3n.awsAccessKeyId : YOUR_ACCESS_KEY, fs.s3n.awsSecretAccessKey : YOUR_SECRET_KEY, fs.s3n.impl : org.apache.hadoop.fs.s3native.NativeS3FileSystem, io.compression.codecs : org.apache.hadoop.io.compress.GzipCodec,org.apache.hadoop.io.compress.DefaultCodec,org.apache.hadoop.io.compress.BZip2Codec,org.apache.hadoop.io.compress.SnappyCodec }注意此方式使用 Hadoop 内置 S3 文件系统而非 Amazon EMRFS不兼容 S3 加密、一致视图等 Amazon 专属特性如需这些特性必须通过下文 使用其他 Hadoop 发行版 的机制把 Amazon EMR Hadoop JAR 提供给 Druid。6.3 安全KerberosHadoop 集群默认情况下 Druid 直接使用本地 Kerberos key cache 中已有的 TGT ticket。但 TGT 有生命周期限制需周期性调用kinit维持有效性为避免额外的外部 cron 脚本可在配置中提供 principal 与 keytab 路径让 Druid 在启动和 job 提交时透明完成认证。其配置对象对应源码 HadoopKerberosConfig.java字段为principal与keytab。PropertyPossible ValuesDescriptionDefaultdruid.hadoop.security.kerberos.principaldruidEXAMPLE.COMPrincipal 用户名emptydruid.hadoop.security.kerberos.keytab/etc/security/keytabs/druid.headlessUser.keytabkeytab 文件路径empty七、其他 Hadoop 发行版与常见问题Druid 开箱即可配合多种 Hadoop 发行版运行。若 Druid 与你使用的 Hadoop 版本之间存在依赖冲突可在社区用户组搜索解决方案或阅读 other-hadoop.md 的 Different Hadoop Versions 文档该文档覆盖了依赖冲突处理、mapreduce.job.classpath等关键配置的实战细节。八、备选方案命令行 Hadoop Indexer 与 IndexTask8.1 命令行 Hadoop Indexer如果不希望为使用 Hadoop 摄入而部署完整的索引服务可运行独立的命令行 Hadoop indexer把 spec 作为命令行参数直接执行详情见 command-line-hadoop-indexer.md。命令行模式下segmentOutputPath与metadataUpdateSpec由用户显式指定这正是 HadoopIOConfig.java 中这两个字段存在的场景并由MetadataStorageUpdaterJob在任务成功后把 segment 列表写入元数据存储。8.2 IndexTask 批式摄入若批式摄入不希望依赖 Hadoop可使用index任务直接在 Druid 进程内完成摄入。相比 Hadoop 方式它更慢、扩展性更低适合小数据量场景详见 tasks.md。选择依据可归纳为数据规模大、已有 Hadoop 集群 →index_hadoop数据量小、无 Hadoop 依赖 →index需要更精细控制 Hadoop job 且不想运行 Overlord → 命令行 Hadoop Indexer。九、小结本文围绕index_hadoop任务给出了从任务 JSON、DataSchema/IOConfig/TuningConfig 到分区规范、集群适配远程集群、EMR、Kerberos的完整实战指南并结合 HadoopIndexTask.java、HadoopTuningConfig.java、StaticPathSpec.java、GranularityPathSpec.java 等源码验证了默认值、执行流程与输入发现机制。结合仓库中的 wikiticker-index.json 示例任务与 batch-ingestion.md 原文档读者可以据此构造自己的第一批 Hadoop 批式摄入任务。赞分享数据库数据分析OLAP大数据实时分析数据仓库后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid7/druid点击查看免费下载相关推荐Apache Druid 基于 Hadoop 的批量摄入Hadoop-based Ingestion完整实战指南Apache Druid 基于 Hadoop 的批量摄入Hadoop based Ingestion完整实战指南 本文是 Apache Druid 官方参考数据库OLAP大数据后端Apache Druid 数据摄入Ingestion完全指南流式与批量摄入方法详解Apache Druid 数据摄入Ingestion完全指南流式与批量摄入方法详解 本篇指南围绕 Apache Druid 的数据摄入Ingestion数据库OLAP大数据后端用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务 Lucky 是一款面向软硬路由的公网管理工具集成端口转发、动态域名DDNS、反向代理、网络后端网络通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表