拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Apache Pulsar 分层存储(Tiered Storage)完全指南:从 Segment 架构到 S3/GCS/文件系统卸载实战

消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载分层存储Tiered Storage是 Apache Pulsar 中一项核心的数据生命周期管理能力它利用 Pulsar 的 Segment 导向架构把 topic backlog 中较旧的消息从 BookKeeper 自动或手动迁移到 S3、Google Cloud StorageGCS或 Hadoop 文件系统等低成本长期存储中同时让生产者和消费者完全无感知。本文将从架构原理出发完整讲解卸载驱动配置、三种后端aws-s3 / google-cloud-storage / filesystem的详细参数、自动与手动卸载操作并结合当前仓库源码tiered-storage 模块与 conf/broker.conf给出可复制、可运行的实战配置。为什么需要分层存储backlog 无上限但成本有上限Pulsar 的 Segment 导向架构segment oriented architecture使得一个 topic 的 backlog 可以近乎无限地增长。在 概念文档 中明确说明backlog 的增长本身几乎没有上限但长时间保留全部数据会带来高昂的存储成本。分层存储正是为了缓解这一成本问题而生backlog 中较老的消息可以被从 BookKeeper 迁移到更便宜的存储介质同时客户端访问 backlog 的体验保持不变——仿佛数据从未被移动过。从成本构成上看写入 BookKeeper 的数据默认会被复制到 3 台物理机器。而一旦 BookKeeper 中的某个 segment 被密封sealed它便成为不可变数据可以被完整拷贝到长期存储。长期存储可以通过 Reed-Solomon 纠删码等机制用更少的物理副本获得同等甚至更强的数据冗余从而显著降低成本。适用场景什么时候该用分层存储分层存储最适合希望长时间保留很长的 backlog的场景。典型例子一个存放用户行为的 topic用于训练推荐系统——当你修改推荐算法时可能需要基于全部历史数据重新训练因此希望这些数据保留数月甚至数年。分层存储的三个关键设计点底层存储不可变只向 log 的最后一个 segment 写入之前的 segment 全部密封、不可变——这是卸载机制能够安全拷贝数据的前提。透明访问数据迁移到长期存储后客户端读写不受影响你依然可以用 Pulsar SQL 查询已卸载 ledger 中的数据。可配置的删除延迟原数据从 BookKeeper 删除前会保留一段配置时间默认 4 小时保证卸载过程出现问题时可以回退。卸载机制原理从 managed ledger 到长期存储一个 Pulsar topic 在底层由一个称为 managed ledger 的 log 承载这个 log 由一组有序的 segment 组成。Pulsar 只向 log 的最后一个 segment 写入数据所有更早的 segment 都被密封其中的数据不可变。分层存储卸载机制正是利用了这一架构当卸载被触发时log 中的 segment 被逐个拷贝到长期存储。除当前正在写入的 segment 外log 中的所有 segment 都可以被卸载。这一机制在源码中的落地可以从两个层面印证驱动抽象managedLedgerOffloadDriver是选择后端驱动的核心配置键见 TieredStorageConfiguration.java。配置类同时支持通用前缀managedLedgerOffload与各供应商的向后兼容前缀如s3ManagedLedgerOffload、gcsManagedLedgerOffload并内置了默认值最大块 64MB、最小块 5MB、读缓冲 1MB、写缓冲 10MB。多部分上传Pulsar 使用多部分对象multipart objects上传 segment 数据。上传过程中 broker 可能崩溃因此文档建议为 bucket 配置生命周期规则让未完成的多部分上传在 1~2 天后自动过期避免为不完整的上传付费。底层 jclouds 覆盖参数可见于 getOverrides 方法并行上传度固定为 1、parts 大小取maxBlockSizeInBytes、SO 超时 25 秒、最大重试 100 次。卸载前的必要前提broker 上必须配置云存储服务的 bucket 和凭证配置的 bucket 必须预先存在否则卸载操作会直接失败卸载触发时你传入希望保留在 BookKeeper 上的 backlog 数据量字节数broker 会从 topic backlog 起始处开始拷贝 segment直到满足该条件为止。配置卸载驱动broker.conf 全参数解析卸载配置全部集中在broker.conf中。最低限度需要配置驱动类型driver、bucket、认证凭证此外还有 region、最大块大小等可选旋钮。支持的驱动类型驱动名对应存储aws-s3Amazon Simple Cloud Storage Servicegoogle-cloud-storageGoogle Cloud StorageGCSfilesystemHadoop 文件系统驱动名大小写不敏感。此外还有一个与aws-s3完全相同的第三种 S3 驱动s3但它要求通过s3ManagedLedgerOffloadServiceEndpoint指定 endpoint URL——适用于 AWS 之外的 S3 兼容数据存储如自建 MinIO、Ceph RGW 等。选择驱动的最小配置managedLedgerOffloadDriveraws-s3在当前仓库的 broker.conf 中与卸载相关的完整配置段还包括# The directory for all the offloader implementations offloadersDirectory./offloaders # Driver to use to offload old data to long term storage # (Possible values: aws-s3, google-cloud-storage, azureblob, aliyun-oss, filesystem) managedLedgerOffloadDriver # Maximum number of thread pool threads for ledger offloading managedLedgerOffloadMaxThreads2 # Maximum prefetch rounds for ledger reading for offloading managedLedgerOffloadPrefetchRounds1仓库实际版本支持的驱动值还包括azureblob、aliyun-oss比 2.3.2 文档列出的三类更丰富说明该能力在持续演进。通用卸载参数删除延迟与自动触发阈值以下两个参数直接控制卸载的数据生命周期行为broker.conf# Delay between a ledger being successfully offloaded to long term storage # and the ledger being deleted from bookkeeper (default is 4 hours) managedLedgerOffloadDeletionLagMs14400000 # The number of bytes before triggering automatic offload to long term storage # (default is -1, which is disabled) managedLedgerOffloadAutoTriggerSizeThresholdBytes-1managedLedgerOffloadDeletionLagMsledger 成功卸载后延迟多久从 BookKeeper 中删除默认 14400000 毫秒4 小时。managedLedgerOffloadAutoTriggerSizeThresholdBytes触发自动卸载的数据量阈值默认 -1 表示禁用实际自动卸载通常由 namespace 策略的 offload threshold 控制见下文。aws-s3 驱动配置Bucket 与 RegionBucket 是存放数据的基本容器Cloud Storage 中所有内容都必须存放在某个 bucket 中。可以用 bucket 组织数据、控制访问权限但与目录/文件夹不同bucket 不能嵌套。s3ManagedLedgerOffloadBucketpulsar-topic-offloadRegion 是 bucket 所在的区域非必填但强烈建议配置不配置时使用默认区域。对 AWS S3 而言默认区域是US East (N. Virginia)。s3ManagedLedgerOffloadRegioneu-west-3与 AWS 的认证方式Pulsar 不直接提供 AWS S3 的认证配置入口而是依赖 AWS Java SDK 的DefaultAWSCredentialsProviderChain。在 AWS IAM 控制台创建好凭证后可以通过以下任一方式配置EC2 实例元数据凭证如果 broker 运行在带有实例配置文件instance profile的 AWS 实例上且没有提供其他机制Pulsar 会自动使用这些凭证。在conf/pulsar_env.sh中设置环境变量注意export关键字必不可少这样才能将变量传入派生进程的环境中export AWS_ACCESS_KEY_IDABC123456789 export AWS_SECRET_ACCESS_KEYded7db27a4558e2ea8bbf0bf37ae0e8521618f366c通过 Java 系统属性把aws.accessKeyId和aws.secretKey加入conf/pulsar_env.sh的PULSAR_EXTRA_OPTSPULSAR_EXTRA_OPTS${PULSAR_EXTRA_OPTS} ${PULSAR_MEM} ${PULSAR_GC} -Daws.accessKeyIdABC123456789 -Daws.secretKeyded7db27a4558e2ea8bbf0bf37ae0e8521618f366c -Dio.netty.leakDetectionLeveldisabled -Dio.netty.recycler.maxCapacity.default1000 -Dio.netty.recycler.linkCapacity1024在~/.aws/credentials中写入访问凭证[default] aws_access_key_idABC123456789 aws_secret_access_keyded7db27a4558e2ea8bbf0bf37ae0e8521618f366c代入 IAM 角色在broker.conf中指定角色 ARN 与会话名通过DefaultAWSCredentialsProviderChain完成角色代入s3ManagedLedgerOffloadRoleaws role arn s3ManagedLedgerOffloadRoleSessionNamepulsar-s3-offload通过pulsar_env.sh指定的凭证需要重启 broker才会生效。在源码层S3 相关凭证与角色配置的键名定义在 TieredStorageConfiguration.javas3ManagedLedgerOffloadCredentialId、s3ManagedLedgerOffloadCredentialSecret、s3ManagedLedgerOffloadRole、s3ManagedLedgerOffloadRoleSessionName。读写块大小配置Pulsar 还提供两个控制 S3 请求尺寸的旋钮broker.confs3ManagedLedgerOffloadMaxBlockSizeInBytes多部分上传中单个 part 的最大尺寸不能小于 5MB默认 64MB67108864。s3ManagedLedgerOffloadReadBufferSizeInBytes从 S3 读回数据时单次读取的块大小默认 1MB1048576。s3ManagedLedgerOffloadMaxBlockSizeInBytes67108864 s3ManagedLedgerOffloadReadBufferSizeInBytes1048576官方建议这两个参数除非确知自己在做什么否则不要改动。google-cloud-storage 驱动配置Bucket 与 Region与 S3 类似GCS 使用 bucket 作为基本容器同样不能嵌套gcsManagedLedgerOffloadBucketpulsar-topic-offloadRegion 非必填但推荐配置。GCS 的 bucket 默认创建在us multi-regional location未配置 region 时使用该默认位置gcsManagedLedgerOffloadRegioneurope-west3与 GCS 的认证方式管理员需要在broker.conf中配置gcsManagedLedgerOffloadServiceAccountKeyFilebroker 才能访问 GCS 服务。该配置指向一个 JSON 文件内含某个服务账号service account的 GCS 凭证。生成服务账号凭证的步骤打开 Google Cloud Console 的 Service accounts 页面选择一个项目或新建一个点击Create service account在弹出的窗口中输入服务账号名称勾选Furnish a new private key若需要授予 G Suite 全域委派权限同时勾选Enable G Suite Domain-wide Delegation点击Create。注意务必确保创建的服务账号拥有操作 GCS 的权限——需要为服务账号分配Storage Admin权限。在broker.conf中指定密钥文件路径gcsManagedLedgerOffloadServiceAccountKeyFile/Users/hello/Downloads/project-804d5e6a6f33.json该键名在源码中的定义为gcsManagedLedgerOffloadServiceAccountKeyFile见 TieredStorageConfiguration.java。另外 broker.conf 的注释还提醒使用 GCS 时需确保项目同时启用了Google Cloud Storage与Google Cloud Storage JSON API在 Developers Console → Apiauth → APIs 中检查。读写块大小配置GCS 驱动的两个旋钮与 S3 完全对应gcsManagedLedgerOffloadMaxBlockSizeInBytes多部分上传单个 part 最大尺寸不能小于 5MB默认 64MB。gcsManagedLedgerOffloadReadBufferSizeInBytes单次读取块大小默认 1MB。gcsManagedLedgerOffloadMaxBlockSizeInBytes67108864 gcsManagedLedgerOffloadReadBufferSizeInBytes1048576同样地非必要不要改动这两个参数。filesystem 驱动配置filesystem 驱动基于 Apache Hadoop用于把卸载数据写入本地或远程文件系统如 HDFS。配置连接地址在broker.conf中配置文件系统地址fileSystemURIhdfs://127.0.0.1:9000配置 Hadoop profile 路径配置文件存放在 Hadoop profile 路径中包含 base path、认证等各类设置fileSystemProfilePath../conf/filesystem_offload_core_site.xml该文件对应仓库中的 conf/filesystem_offload_core_site.xml其默认内容为configuration !--file system uri, necessary-- property namefs.defaultFS/name value/value /property property namehadoop.tmp.dir/name valuepulsar/value /property property nameio.file.buffer.size/name value4096/value /property property nameio.seqfile.compress.blocksize/name value1000000/value /property property nameio.seqfile.compression.type/name valueBLOCK/value /property property nameio.map.index.interval/name value128/value /property /configurationtopic 数据的存储模型使用org.apache.hadoop.io.MapFile因此可以复用 HadoopMapFile的全部配置项如块压缩、索引间隔等。filesystem 驱动对应的实现类为 FileSystemManagedLedgerOffloader.java其工厂类位于 FileSystemLedgerOffloaderFactory.java。配置自动卸载namespace 级 offload thresholdnamespace 策略可以配置为在达到阈值时自动卸载数据。阈值基于该 topic 在 Pulsar 集群上已存储的数据量topic 达到阈值后即触发一次卸载操作。将阈值设为负值可禁用自动卸载设为 0 则 broker 会尽可能早地执行卸载。$ bin/pulsar-admin namespaces set-offload-threshold --size 10M my-tenant/my-namespace自动卸载是在新的 segment 加入 topic log 时触发的。如果设置了阈值但 topic 几乎不产生消息那么在当前 segment 写满之前卸载不会发生。从源码看namespace 级 offload threshold 最终写入OffloadPoliciesImpl的managedLedgerOffloadThresholdInBytes字段相关逻辑位于 NamespacesBase.javabroker 在启动时也会通过OffloadPoliciesImpl.create(...)将全局配置转换为默认卸载策略见 PulsarService.java。配置已卸载消息的读取优先级默认情况下消息被卸载到长期存储后broker 会从长期存储读取这些消息但消息在 BookKeeper 中仍会保留一段时间时长由管理员配置默认 4 小时。对于同时存在于 BookKeeper 与长期存储中的消息如果希望优先从 BookKeeper 读取可以通过以下命令修改读取优先级# default value for -orp is tiered-storage-first $ bin/pulsar-admin namespaces set-offload-policies my-tenant/my-namespace -orp bookkeeper-first $ bin/pulsar-admin topics set-offload-policies my-tenant/my-namespace/topic1 -orp bookkeeper-first-orp参数的默认值为tiered-storage-first即优先从分层存储读取可改为bookkeeper-first让 broker 优先读取 BookKeeper 中的数据。在 PersistentTopicsBase.java 中可以看到topic 级策略会与 namespace 级策略通过OffloadPoliciesImpl.mergeConfiguration合并实现topic 覆盖 namespace的优先级体系。手动触发卸载与状态查询卸载也可以通过 broker 上的 REST 端点手动触发pulsar-admin提供了调用该端点的 CLI。手动触发时必须指定将保留在 BookKeeper 本地的 backlog 最大字节数size-threshold。卸载机制会从 topic backlog 的起始处开始卸载 segment直到满足该条件。$ bin/pulsar-admin topics offload --size-threshold 10M my-tenant/my-namespace/topic1 Offload triggered for persistent://my-tenant/my-namespace/topic1 for messages before 2:0:-1offload命令不会等待卸载完成即返回。要查看卸载状态使用offload-status$ bin/pulsar-admin topics offload-status my-tenant/my-namespace/topic1 Offload is currently running如需等待卸载完成加-w标志$ bin/pulsar-admin topics offload-status -w my-tenant/my-namespace/topic1 Offload was a success卸载失败时的排错如果卸载出错错误信息会传播到offload-status命令。例如下面的输出展示了一次因 S3 认证失败导致的卸载错误$ bin/pulsar-admin topics offload-status persistent://public/default/topic1 Error in offload null Reason: Error offloading: org.apache.bookkeeper.mledger.ManagedLedgerException: java.util.concurrent.CompletionException: com.amazonaws.services.s3.model.AmazonS3Exception: Anonymous users cannot initiate multipart uploads. Please authenticate. (Service: Amazon S3; Status Code: 403; Error Code: AccessDenied; Request ID: 798758DE3F1776DF; S3 Extended Request ID: dhBFz/lZm1oiG/oBEepeNlhrtsDlzoOhocuYMpKihQGXe6EG8puRGOkK6UwqzVrMXTWBxxHcSg), S3 Extended Request ID: dhBFz/lZm1oiG/oBEepeNlhrtsDlzoOhocuYMpKihQGXe6EG8puRGOkK6UwqzVrMXTWBxxHcSg这个例子很好地印证了前文提到的 S3 认证链报错Anonymous users cannot initiate multipart uploadsStatus Code 403Error Code AccessDenied说明 broker 未能通过DefaultAWSCredentialsProviderChain取得任何有效凭证请按上文五种认证方式之一完成配置。用 Pulsar SQL 查询已卸载的数据ledger 被卸载到长期存储后你仍然可以查询已卸载 ledger 中的数据——Pulsar SQL 支持对分层存储中的数据执行查询。这使得归档数据与在线数据在查询层面保持统一离线分析任务可以直接作用于长期存储中的数据无需先导回 BookKeeper。小结分层存储配置清单围绕本仓库的文档与源码一个完整的分层存储启用流程可以归纳为在broker.conf中设置managedLedgerOffloadDriveraws-s3/google-cloud-storage/filesystem/s3等并确保offloadersDirectory指向已部署相应 offloader NAR 包的目录配置对应驱动的 buckets3ManagedLedgerOffloadBucket/gcsManagedLedgerOffloadBucket与 region确保 bucket 预先存在配置认证凭证S3 走DefaultAWSCredentialsProviderChain环境变量、系统属性、credentials 文件、实例元数据或 IAM 角色代入GCS 指定gcsManagedLedgerOffloadServiceAccountKeyFilefilesystem 配置fileSystemURI与fileSystemProfilePath可选按需调整块大小参数与managedLedgerOffloadDeletionLagMs默认 4 小时删除延迟设置自动卸载bin/pulsar-admin namespaces set-offload-threshold --size 10M my-tenant/my-namespace或手动执行bin/pulsar-admin topics offload --size-threshold ...并用offload-status -w跟踪结果按需用set-offload-policies -orp bookkeeper-first调整已卸载消息的读取优先级。相关文档与源码入口分层存储概念文档分层存储实战 cookbookbroker 卸载配置段Hadoop filesystem 卸载配置模板jcloud 卸载配置实现filesystem 卸载实现赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 分层存储Tiered Storage实战指南卸载旧数据到 S3 / GCS / 文件系统Apache Pulsar 分层存储Tiered Storage实战指南卸载旧数据到 S3 / GCS / 文件系统 分层存储Tiered Storag消息队列后端流处理Apache Pulsar 分层存储Tiered Storage完全指南从架构原理到 S3/GCS/文件系统 Offload 实战Apache Pulsar 分层存储Tiered Storage完全指南从架构原理到 S3/GCS/文件系统 Offload 实战 Pulsar 的分层存消息队列后端流处理Apache Pulsar 分层存储Tiered Storage完全指南架构原理与 S3/GCS/文件系统 Offloader 实战Apache Pulsar 分层存储Tiered Storage完全指南架构原理与 S3/GCS/文件系统 Offloader 实战 分层存储Tiere消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门