尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Pulsar 分层存储(Tiered Storage)实战指南:卸载旧数据到 S3 / GCS / 文件系统

Apache Pulsar 分层存储(Tiered Storage)实战指南:卸载旧数据到 S3 / GCS / 文件系统 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载分层存储Tiered Storage是 Apache Pulsar 面向“超长保留期”数据场景的核心能力它允许把主题Topic日志中的旧数据段Segment从 BookKeeper 卸载Offload到 S3、GCS 或 Hadoop 文件系统等长期存储中从而释放 BookKeeper 磁盘空间、降低存储成本同时保证历史数据仍可被消费与查询。本文以 Apache Pulsar 官方 Cookbookversion-2.3.2 版为骨架结合当前仓库源码tiered-storage/jcloud、tiered-storage/file-system、conf/broker.conf 等展开带你完整掌握分层存储的触发机制、三种 offload 驱动aws-s3 / google-cloud-storage / filesystem的配置方法、自动与手动卸载的运维命令以及异常排查要点。何时该使用分层存储分层存储适合“某个主题需要保留非常长的 backlog 很长时间”的场景。典型例子主题中保存着用户行为数据用来训练推荐系统——你可能希望长期保留全部历史以便日后调整推荐算法时能基于完整历史重新训练。如果这些数据始终留在 BookKeeper 内存储成本会持续攀升通过分层存储把早期数据段搬移到低成本对象存储既保留了完整历史又控制了热存储开销。值得强调的是卸载到长期存储后这些 ledger 数据仍然可以通过 Pulsar SQL 查询数据的可消费性、可查询性不会因卸载而丢失。卸载机制Offloading MechanismPulsar 中一个主题由一条被称为managed ledger的日志承载这条日志由一组**有序的分段segment**组成。Pulsar 只会向日志的最后一个分段写入数据此前的所有分段都被sealed封存其中的数据不可变——这就是 Pulsar 的分段导向架构segment oriented architecture。分层存储正是利用了这一架构当卸载被触发时日志中的分段会被逐段复制到长期存储。除当前正在写入的分段外日志中所有其他分段都可以被卸载。几个关键运维要点管理员必须在 broker 上配置云存储服务的bucket桶与认证凭据。bucket 必须预先创建。如果 bucket 不存在卸载操作会直接失败。Pulsar 使用multipart upload分片上传来上传分段数据上传过程中 broker 有可能崩溃导致产生不完整的上传任务。官方建议在 bucket 上添加一条生命周期规则让不完整的分片上传在一到两天后过期避免为未完成的上传持续计费。从源码看卸载的核心实现在 BlobStoreManagedLedgerOffloader.java 及其对应的测试 BlobStoreManagedLedgerOffloaderTest.java它通过 jclouds 的BlobStore接口完成对象存储交互。底层配置类TieredStorageConfiguration所有分层存储配置最终由 TieredStorageConfiguration.java 承载。该类从全局broker.conf中读取配置主要字段包括驱动键managedLedgerOffloadDriverBLOB_STORE_PROVIDER_KEY元数据字段bucket、region、serviceEndpoint、maxBlockSizeInBytes、readBufferSizeInBytes 等向后兼容键S3 使用s3ManagedLedgerOffload*前缀GCS 使用gcsManagedLedgerOffload*前缀见getBackwardCompatibleKey默认值最大分块大小默认64MBgetMaxBlockSizeInBytes返回64 * MB最小分块大小5MB读缓冲区默认1MB写缓冲区默认10MB最大分段大小默认1GB。其中validate()方法会校验region 与 serviceEndpoint 至少指定一个、bucket 不能为空、MaxBlockSizeInBytes不能小于 5MB否则抛出IllegalArgumentException。配置卸载驱动Configuring the Offload Driver卸载功能在 conf/broker.conf 中配置对应小节### --- Ledger Offloading --- ###位于文件约第 1360 行起。至少需要配置driver驱动、bucket桶、认证凭据。此外还有 region、块大小等可调参数。当前仓库支持的驱动类型有aws-s3AWS 简单云存储服务S3google-cloud-storageGoogle 云存储GCSfilesystem基于 Apache Hadoop 的文件系统存储。驱动名称不区分大小写。另外还有第三种 S3 兼容驱动s3它与aws-s3行为一致但必须通过s3ManagedLedgerOffloadServiceEndpoint指定 endpoint URL适用于对接非 AWS 的 S3 兼容存储例如自建 MinIO、Ceph RGW、阿里云 OSS 等。最小配置示例managedLedgerOffloadDriveraws-s3在 broker.conf 中与卸载相关的全局参数还有# 所有 offloader 实现所在的目录 offloadersDirectory./offloaders # 用于卸载旧数据到长期存储的驱动可选值aws-s3, google-cloud-storage, azureblob, aliyun-oss, filesystem managedLedgerOffloadDriver # ledger 卸载的线程池最大线程数 managedLedgerOffloadMaxThreads2 # ledger 读取用于卸载的最大预取轮数 managedLedgerOffloadPrefetchRounds1 # 一个 ledger 成功卸载到长期存储后本地副本延迟删除的时间默认 14400000 ms 4 小时 managedLedgerOffloadDeletionLagMs14400000 # 触发自动卸载到长期存储的字节数阈值默认 -1表示关闭自动卸载 managedLedgerOffloadAutoTriggerSizeThresholdBytes-1从 broker.conf 第 1365 行的注释可以看出除了文档重点介绍的三种驱动当前仓库还支持azureblobAzure BlobStore通过环境变量AZURE_STORAGE_ACCOUNT/AZURE_STORAGE_ACCESS_KEY认证与aliyun-oss阿里云 OSS兼容 S3 API通过ACCESS_KEY_ID/ACCESS_KEY_SECRET环境变量认证二者均在 JCloudBlobStoreProvider.java 中以枚举形式实现。aws-s3 驱动配置Bucket 与 RegionBucket 是存储数据的基本容器云存储中的所有数据都必须存放在某个 bucket 中。你可以用 bucket 组织数据、控制访问权限但与目录/文件夹不同bucket 不能嵌套。s3ManagedLedgerOffloadBucketpulsar-topic-offloadRegion 是 bucket 所在的区域。Region 不是必填项但强烈建议配置若不配置将使用默认区域。对 AWS S3 而言默认区域是US East (N. Virginia)。AWS 官方“Regions and Endpoints”页面提供了更完整的区域信息。s3ManagedLedgerOffloadRegioneu-west-3从源码看TieredStorageConfiguration对 region 与 endpoint 的校验逻辑是“二者至少指定其一”见JCloudBlobStoreProvider.VALIDATIONif (Strings.isNullOrEmpty(config.getRegion()) Strings.isNullOrEmpty(config.getServiceEndpoint())) { throw new IllegalArgumentException( Either Region or ServiceEndpoint must specified for config.getDriver() offload); }AWS 认证方式Pulsar 不提供针对 AWS S3 的直接认证配置入口而是依赖 AWS SDK 的DefaultAWSCredentialsProviderChain。在 AWS IAM 控制台创建凭据后可以通过以下方式之一配置使用 EC2 实例元数据凭据如果你的 broker 运行在带实例配置文件instance profile的 AWS 实例上且没有提供其他机制Pulsar 会自动使用这些凭据。在conf/pulsar_env.sh中设置环境变量AWS_ACCESS_KEY_ID与AWS_SECRET_ACCESS_KEYexport AWS_ACCESS_KEY_IDABC123456789 export AWS_SECRET_ACCESS_KEYded7db27a4558e2ea8bbf0bf37ae0e8521618f366cexport很重要这样才能让变量对派生的子进程环境可见。在conf/pulsar_env.sh的PULSAR_EXTRA_OPTS中添加 Java 系统属性aws.accessKeyId与aws.secretKeyPULSAR_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_keyded7db27a4558e2ea8bbf0bf37ae0e8521618f366cAssume 一个 IAM 角色角色切换在broker.conf中指定如下配置即可通过DefaultAWSCredentialsProviderChain完成角色切换s3ManagedLedgerOffloadRoleaws role arn s3ManagedLedgerOffloadRoleSessionNamepulsar-s3-offload注意在 pulsar_env 中指定的凭据需要重启 broker才会生效。从 JCloudBlobStoreProvider.java 的AWS_CREDENTIAL_BUILDER可以看到三种情况的完整分支若配置了s3ManagedLedgerOffloadCredentialId/s3ManagedLedgerOffloadCredentialSecret即源码中定义的S3_ID_FIELD/S3_SECRET_FIELD则用静态凭据否则若未配置 role则走DefaultAWSCredentialsProviderChain否则走STSAssumeRoleSessionCredentialsProvider。源码注释特别强调凭据构建被延迟到真正需要时才执行以支持 EC2 元数据这类可能过期的 session token。读写块大小配置Pulsar 提供两个与 AWS S3 请求粒度相关的参数s3ManagedLedgerOffloadMaxBlockSizeInBytesmultipart upload 中单个 “part” 的最大大小不能小于 5MB默认 64MBs3ManagedLedgerOffloadReadBufferSizeInBytes从 S3 读回数据时每次单独读取的块大小默认 1MB。两个参数除非清楚自己在做什么否则不要随意改动。对应到 broker.conf 中的默认值分别为6710886464MB与10485761MB与源码中getMaxBlockSizeInBytes()返回64 * MB、getReadBufferSizeInBytes()返回MB的默认逻辑一致。google-cloud-storage 驱动配置与 S3 类似GCS 的 bucket 同样是数据的容器不可嵌套gcsManagedLedgerOffloadBucketpulsar-topic-offloadRegion 同样是推荐配置。对 GCS 而言bucket 默认创建在us多区域位置us multi-regional locationgcsManagedLedgerOffloadRegioneurope-west3GCS 认证方式管理员需要在broker.conf中配置gcsManagedLedgerOffloadServiceAccountKeyFile它是一个包含服务账号凭据的 JSON 文件。生成服务账号凭据或查看已有凭据的步骤打开 Google Cloud 控制台的 Service accounts 页面选择一个项目或新建一个点击Create service account创建服务账号在创建窗口中输入服务账号名称勾选Furnish a new private key生成新的私钥。如需授予 G Suite 域级委派权限可同时勾选Enable G Suite Domain-wide Delegation点击Create。注意确保所创建的服务账号拥有操作 GCS 的权限——需要在 IAM 中为服务账号授予Storage Admin权限。配置示例gcsManagedLedgerOffloadServiceAccountKeyFile/Users/hello/Downloads/project-804d5e6a6f33.json源码实现GOOGLE_CLOUD_STORAGE.buildCredentials会读取该 JSON 文件内容并通过GoogleCredentialsFromJson构造 jclouds 凭据String gcsKeyContent Files.asCharSource( new File(config.getConfigProperty(GCS_ACCOUNT_KEY_FILE_FIELD)), Charset.defaultCharset()).read(); config.setProviderCredentials(() - new GoogleCredentialsFromJson(gcsKeyContent).get());若文件读取失败例如路径错误会抛出IllegalArgumentException并记录Cannot read GCS service account credentials file错误日志。读写块大小配置gcsManagedLedgerOffloadMaxBlockSizeInBytesmultipart upload 中单个 part 的最大大小不能小于 5MB默认 64MBgcsManagedLedgerOffloadReadBufferSizeInBytes从 GCS 读回数据时每次单独读取的块大小默认 1MB。同样地默认值不建议随意改动。这两个配置在 broker.conf 中默认分别为67108864与1048576。filesystem 驱动配置filesystem 驱动基于 Apache Hadoop将数据卸载到文件系统如 HDFS。卸载数据使用org.apache.hadoop.io.MapFile模型存储因此 Hadoop 中MapFile支持的所有配置项在这里都可以使用。配置连接地址在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.xmlconfiguration !--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 /configuration其中各配置项含义fs.defaultFS必填的文件系统 URI与broker.conf中fileSystemURI保持一致hadoop.tmp.dirHadoop 临时目录io.file.buffer.size文件读写缓冲区大小io.seqfile.compress.blocksizeSequenceFile 压缩块大小io.seqfile.compression.type压缩类型示例为BLOCKio.map.index.intervalMapFile 索引写入间隔。filesystem 驱动的核心实现在 FileSystemManagedLedgerOffloader.java 与 FileStoreBackedReadHandleImpl.java读回句柄入口为 FileSystemLedgerOffloaderFactory.java工厂机制与 jcloud 的JCloudLedgerOffloaderFactory相对应。配置自动卸载通过命名空间Namespace策略可以在达到阈值后自动触发卸载。阈值基于主题在 Pulsar 集群上已存储的数据量一旦主题达到阈值就会触发一次卸载操作。将阈值设为负值会禁用自动卸载将阈值设为0会让 broker 尽可能快地卸载数据。$ bin/pulsar-admin namespaces set-offload-threshold --size 10M my-tenant/my-namespace自动卸载只在向主题日志新增一个分段时运行。如果你对某个命名空间设置了阈值但主题几乎没有新消息产生那么卸载要等到当前分段写满后才会发生。对应 CLI 的实现位于 CmdNamespaces.java 的SetOffloadThreshold命令其参数说明明确指出“Negative values disable automatic offload. 0 triggers offloading as soon as possible”与文档表述完全一致。后端 REST 接口则实现在 Namespaces.java 的setOffloadThresholdv2 路径为/{tenant}/{namespace}/offloadThreshold。配置卸载后的读取优先级默认情况下消息一旦卸载到长期存储broker 会从长期存储读取但消息在 BookKeeper 中还会保留一段时间保留时长取决于管理员配置。对于同时存在于 BookKeeper 与长期存储的消息如果希望优先从 BookKeeper 读取可以通过命令修改这一配置# 默认值 -orp 为 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即--offloadReadPriorityOffloadedReadPriority相关策略实现在 OffloadPoliciesImplorg.apache.pulsar.common.policies.data.OffloadPoliciesImpl中可对命名空间级与主题级分别设置。此外命名空间级别还提供offloadDeletionLag设置——即 ledger 分段卸载成功后本地 BookKeeper 副本延迟删除的时长CmdNamespaces.java 中的SetOffloadDeletionLag命令对应 broker.conf 中的全局默认值managedLedgerOffloadDeletionLagMs144000004 小时。手动触发卸载卸载也可以手动触发。Pulsar broker 提供对应的 REST 端点CLI 工具pulsar-admin会帮你调用它。手动触发时必须指定本地 BookKeeper 中保留的最大 backlog 字节数--size-threshold。卸载机制会从主题 backlog 的起始位置开始卸载分段直到满足“本地保留量不超过该值”的条件。$ 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:-1注意触发卸载的命令不会等待卸载完成。要检查卸载状态使用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命令透传出来。例如典型的认证失败场景$ 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上面的错误表明broker 尝试以匿名身份发起 S3 multipart upload 被拒绝HTTP 403 AccessDenied——原因通常是 S3 认证凭据没有正确配置比如只设置了s3ManagedLedgerOffloadBucket而未提供任何凭据可按上文“AWS 认证方式”的任一路径补齐凭据并重启 broker 后重试。小结分层存储把 Pulsar 的分段导向架构优势与云对象存储的成本优势结合起来是管理超长 backlog 的标准方案。实践要点可归纳为提前准备存储桶bucket 必须预先存在否则卸载失败建议为 bucket 配置生命周期规则清理不完整的分片上传驱动与凭据在 conf/broker.conf 中设置managedLedgerOffloadDriver并完成对应驱动认证AWS 走DefaultAWSCredentialsProviderChainGCS 走服务账号 JSON 文件filesystem 走 Hadoopcore-site.xml运维命令set-offload-threshold配置自动卸载阈值topics offload/offload-status手动触发并跟踪状态set-offload-policies -orp调整卸载后的读取优先级错误排查offload-status会透传底层异常如 S3 403 AccessDenied是定位配置问题的第一入口。如果想深入源码推荐从 TieredStorageConfiguration.java 与 JCloudBlobStoreProvider.java 入手它们完整覆盖了各驱动从配置解析、校验到凭据构建、BlobStore 创建的全过程。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 分层存储Tiered Storage实战指南将 BookKeeper 旧数据卸载到 S3 与 GCSApache Pulsar 分层存储Tiered Storage实战指南将 BookKeeper 旧数据卸载到 S3 与 GCS 分层存储Tiered消息队列后端流处理Apache Pulsar 分层存储Tiered Storage实战指南把历史积压数据卸载到 S3、GCS 与 Hadoop 文件系统Apache Pulsar 分层存储Tiered Storage实战指南把历史积压数据卸载到 S3、GCS 与 Hadoop 文件系统 分层存储Tier消息队列后端流处理Apache Pulsar 分层存储Tiered Storage实战指南将历史消息卸载到 AWS S3、GCS 与 Hadoop 文件系统Apache Pulsar 分层存储Tiered Storage实战指南将历史消息卸载到 AWS S3、GCS 与 Hadoop 文件系统 Apache消息队列后端流处理上一篇工业级音频AI新突破Step-Audio 2模型重构语音交互技术边界下一篇AnythingLLM 5 分钟跑起来本地私有知识库的 Docker 部署指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表