ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Druid 的 S3 兼容深度存储与 StaticS3Firehose 配置实战

Apache Druid 的 S3 兼容深度存储与 StaticS3Firehose 配置实战 数据库数据分析OLAP大数据实时分析数据仓库后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid7/druid点击查看免费下载本文基于 Apache Druid本仓库为 druid7/druid版本 0.10.1-SNAPSHOT的官方扩展文档系统讲解如何将 S3 及一切兼容 S3 API 的对象存储如 Google Cloud Storage接入 Druid既可作为集群的深度存储deep storage持久化已索引的 Segment也可借助StaticS3Firehose从预定义的 S3 对象列表批量摄取数据。读完本文你将掌握druid-s3-extensions扩展的加载方式、完整参数配置、底层实现原理以及可复制到生产环境的配置示例。druid-s3-extensions 扩展安装与加载S3 相关的全部能力深度存储、Firehose 摄取、任务日志归档都由druid-s3-extensions这一核心扩展提供。使用前必须先将其加入扩展加载列表官方文档要求先加载扩展再使用。在conf/druid/_common/common.runtime.properties中配置druid.extensions.loadList[druid-s3-extensions]druid-s3-extensions位于 extensions-core/s3-extensions 模块下Maven 坐标为io.druid.extensions:druid-s3-extensions。它是 Druid 随发行包内置的核心扩展之一在扩展列表中对其定位的官方描述是Interfacing with data in AWS S3, and using S3 as deep storage.如果你的环境需要单独拉取扩展可以使用pull-deps工具参见加载扩展文档。加载完成后Druid 进程启动时会通过S3StorageDruidModule完成该扩展内全部组件的注册与装配其核心注册逻辑位于 S3StorageDruidModule.java将druid.s3前缀绑定到AWSCredentialsConfigAWS 凭据将druid.storage前缀绑定到S3DataSegmentPusherConfig与S3DataSegmentArchiverConfigSegment 推送 / 归档配置注册 Segment 的 puller、killer、mover、archiver、pusher、finder 以及任务日志task logs等组件将druid.indexer.logs前缀绑定到S3TaskLogsConfig。S3 兼容深度存储定位与工作原理深度存储是 Druid 中承担已落盘数据持久化角色的存储层。Druid 官方对 S3 深度存储的定义是S3-compatible deep storage 本质上就是 S3 本身或者像 Google Storage 这样暴露了与 S3 相同 API 的服务。因此只要对象存储服务兼容 S3 协议均可通过同一套配置接入。在典型集群中深度存储承载两类职责Segment 持久化索引任务Indexing Task将生成的 Segment 压缩后推送到 S3Historical 节点在需要时再从 S3 拉取加载中间结果/日志归档索引服务的任务日志可写入 S3便于集中查看与排查。druid-s3-extensions使用 JetS3t 客户端库org.jets3t.service.impl.rest.httpclient.RestS3Service访问对象存储由S3StorageDruidModule中的getRestS3Service方法根据凭据类型普通 AccessKey 或 AWS 会话凭据构造对应的客户端实例。深度存储核心配置必需配置参数在common.runtime.properties中设置以下四个参数即官方文档定义的必需配置PropertyPossible ValuesDescriptionDefaultdruid.s3.accessKeyS3 access key访问密钥 ID。Must be set.必须设置druid.s3.secretKeyS3 secret key秘密访问密钥。Must be set.必须设置druid.storage.bucketBucket to store in存储使用的桶名。Must be set.必须设置druid.storage.baseKeyBase key prefix to use, i.e. what directory基础 key 前缀即视为目录的前缀。Must be set.必须设置其中前两个参数对应AWSCredentialsConfig类中的accessKey、secretKey字段见 aws-common/src/main/java/io/druid/common/aws/AWSCredentialsConfig.java后两个参数对应 S3DataSegmentPusherConfig.java 中的bucket、baseKey字段。完整配置示例参考官方集群部署教程中的 S3 章节同时使用 S3 作为深度存储与索引服务日志存储的完整配置如下需注释掉本地存储的对应项druid.extensions.loadList[druid-s3-extensions] #druid.storage.typelocal #druid.storage.storageDirectoryvar/druid/segments druid.storage.types3 druid.storage.bucketyour-bucket druid.storage.baseKeydruid/segments druid.s3.accessKey... druid.s3.secretKey... #druid.indexer.logs.typefile #druid.indexer.logs.directoryvar/druid/indexing-logs druid.indexer.logs.types3 druid.indexer.logs.s3Bucketyour-bucket druid.indexer.logs.s3Prefixdruid/indexing-logs关键点说明druid.storage.types3显式声明深度存储后端为 S3druid.storage.baseKeydruid/segments相当于在桶内指定一个根目录所有 Segment 都会写在该前缀之下若要让索引服务任务日志也走 S3需要同时设置druid.indexer.logs.types3、druid.indexer.logs.s3Bucket与druid.indexer.logs.s3Prefix这些参数也可以以 JVM 系统属性形式传递例如在将段插入数据库场景中使用-Ddruid.s3.accessKey... -Ddruid.s3.secretKey... -Ddruid.storage.bucketyour-bucket -Ddruid.storage.baseKeydruid/storage/wikipedia。从源码看配置如何被消费Segment 推送push逻辑由 S3DataSegmentPusher.java 实现其核心流程如下通过S3Utils.constructSegmentPath(config.getBaseKey(), inSegment)计算目标路径格式为{baseKey}/{dataSource}/{时间区间}/{版本}/{分区号}/index.zip将 Segment 目录整体压缩为临时index.zip以流式方式putObject上传到配置的 bucket同时上传同目录下的descriptor.jsonSegment 描述文件路径为{上述路径}/descriptor.json上传完成后Segment 的loadSpec会被更新为{type: s3_zip, bucket: ..., key: ...}供后续节点据此加载。type: s3_zip对应的加载端实现是 S3LoadSpec.javaHistorical 节点通过内部的S3DataSegmentPuller按 bucket key 从 S3 拉取 Segment 文件。加载/推送等能力由 S3Utils.java 提供路径拼接、对象列举与重试等公共工具方法。可选进阶参数除了官方文档列出的四个必需参数外从S3DataSegmentPusherConfig的源码字段S3DataSegmentPusherConfig.java还可以看到以下可选配置Property含义默认值druid.storage.disableAcl为true时上传 Segment 及其描述文件时不附带bucket-owner-full-control的 canned ACL适用于目标桶策略已控制权限、无需写入 ACL 的场景。falsedruid.storage.maxListingLength单次 S3 List 请求最多返回的对象数量非负整数用于分页枚举桶内对象防止大桶下一次列举拉爆内存。1000maxListingLength的实际用途体现在S3Utils.storageObjectsIterator它通过listObjectsChunked(bucket, prefix, null, maxListingLength, priorLastKey)以游标方式分页迭代桶内对象。另外AWSCredentialsConfig还定义了fileSessionCredentials字段绑定前缀同样为druid.s3用于以文件方式提供 AWS 会话凭据从源码结构看当凭据类型为AWSSessionCredentials时S3StorageDruidModule会使用AWSSessionCredentialsAdapter适配为 JetS3t 可识别的会话凭据。与 Hadoop 批处理任务的衔接S3 作为深度存储时Hadoop 索引任务需要知道数据的 S3 路径。S3DataSegmentPusher.getPathForHadoop()会返回形如s3n://{bucket}/{baseKey}的 Hadoop 兼容路径见 S3DataSegmentPusher.java。这意味着基于 Hadoop 的批量索引任务可以直接以s3n://协议读写该深度存储路径。StaticS3Firehose从 S3 对象列表批量摄取功能定位StaticS3Firehose是一种静态 Firehose数据源它从一个预定义的 S3 对象列表中读取事件读取完成后 Firehose 即干涸结束适用于从若干固定文件批量导入数据的场景。其官方定义为This firehose ingests events from a predefined list of S3 objects.示例 Spec在索引任务Indexing Task的firehose字段中按如下方式声明来自官方文档firehose : { type : static-s3, uris: [s3://foo/bar/file.gz, s3://bar/foo/file2.gz] }参数说明propertydescriptiondefaultrequired?typeThis should be static-s3固定值标识 Firehose 类型N/AyesurisJSON array of URIs where s3 files to be ingested are located待摄取 S3 文件的 URI 数组N/Ayes源码级实现解析StaticS3FirehoseFactory见 extensions-core/s3-extensions/src/main/java/io/druid/firehose/s3/StaticS3FirehoseFactory.java的实现细节值得注意URI 校验构造时逐一校验每个 URI 的 scheme 必须为s3inputURI.getScheme().equals(s3)否则直接抛出IllegalArgumentException。因此uris中不能混入s3n://、http://等其他协议对象解析从 URI 的authority部分解析出 bucket从path部分去掉开头的/解析出 object key按序读取将所有 URI 放入队列通过FileIteratingFirehose逐行迭代每个对象的输入流直至队列为空hasNext()返回 falseFirehose 随即结束gzip 自动解压当 object key 以.gz结尾时自动用CompressionUtils.gzipInputStream包裹输入流进行解压见 StaticS3FirehoseFactory.java无需在配置中额外声明压缩格式。该 Firehose 的序列化/反序列化由测试 StaticS3FirehoseFactoryTest.java 验证测试使用与官方示例完全一致的 URIs3://foo/bar/file.gz、s3://bar/foo/file2.gz进行 JSON 序列化与反序列化往返并断言结果相等证明上述 spec 格式可直接被任务解析器接受。需要说明的是StaticS3FirehoseFactory依赖通过JacksonInject(s3Client)注入的RestS3Service客户端由扩展模块的S3FirehoseDruidModule提供因此使用该 Firehose 前同样必须保证druid-s3-extensions已加载且druid.s3.accessKey/druid.s3.secretKey已正确配置。运维与排障要点上传失败自动重试S3Utils.retryS3Operation为所有关键 S3 写操作Segment 推送、任务日志上传、对象列举等提供了重试封装最多重试 10 次且只对可恢复的失败重试——包括底层 IO 异常IOException以及错误码为RequestTimeout的请求超时对于AccessDenied、NoSuchKey、NoSuchBucket等服务级错误则不会重试而是直接抛出见 S3Utils.java。对象存在性判断的语义S3Utils.isObjectInBucket对getObjectDetails返回的 404NoSuchKey/NoSuchBucket判定为对象不存在而对AccessDenied错误则判定为对象存在但当前凭据无权访问。因此在排查段加载失败时若日志出现 AccessDenied 语义应优先检查 IAM/桶策略权限而不是认为对象缺失。任务日志的读取方式S3TaskLogs.streamTaskLog见 S3TaskLogs.java支持带偏移量的流式读取日志对象的 key 格式为{s3Prefix}/{taskId}/log读取时通过Range请求start/end字节区间实现从任意偏移继续读取方便 Overlord 在任务运行期间持续追读日志。延伸阅读深度存储总体说明docs/content/dependencies/deep-storage.md扩展加载方式docs/content/operations/including-extensions.md通用配置参考含 S3 参数速查docs/content/configuration/index.mdS3 深度存储的集群部署步骤docs/content/tutorials/cluster.md手动向 S3 插入 Segment 的示例docs/content/operations/insert-segment-to-db.md扩展源码extensions-core/s3-extensions/src/main/java/io/druid/storage/s3 与 extensions-core/s3-extensions/src/main/java/io/druid/firehose/s3赞分享数据库数据分析OLAP大数据实时分析数据仓库后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid7/druid点击查看免费下载相关推荐hindex生产部署清单上线前必须确认的8项配置要点hindex生产部署清单上线前必须确认的8项配置要点 hindex 是一个 100% Java 实现的 HBase 二级索引Secondary Index数据库OLAP大数据后端Apache Druid 深度存储迁移实战从本地存储平滑迁往 S3 / HDFS / 新本地路径Apache Druid 深度存储迁移实战从本地存储平滑迁往 S3 / HDFS / 新本地路径 Apache Druid 的 深度存储deep stora数据库OLAP大数据后端Apache Druid 深度存储接入 Apache Cassandradruid-cassandra-storage 扩展配置与实现原理Apache Druid 深度存储接入 Apache Cassandradruid cassandra storage 扩展配置与实现原理 Apache Dr数据库OLAP大数据后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表