ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Strimzi 分层存储系统测试深度解析:基于 Aiven Tiered Storage 插件的 NFS 与 S3 端到端验证

Strimzi 分层存储系统测试深度解析:基于 Aiven Tiered Storage 插件的 NFS 与 S3 端到端验证 Strimzi 分层存储系统测试深度解析基于 Aiven Tiered Storage 插件的 NFS 与 S3 端到端验证【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator分层存储Tiered Storage是 Kafka 3.9.0 起正式可用的特性允许将历史日志段卸载到独立存储系统实现热数据留在本地块存储、冷数据下沉到低成本远程存储的架构。本文基于 Strimzi Kafka Operator 仓库中的系统测试套件 TieredStorageST完整解析 Strimzi 如何通过Kafka自定义资源的tieredStorage配置接入 Aiven Tiered Storage 插件并借助 SeaweedFSS3 兼容对象存储与 NFS 两种后端端到端验证日志段上传、本地删除、远程读取与远程删除的完整生命周期。读完本文你将掌握 Strimzi 分层存储的 CR 配置方法、关键调优参数以及这套系统测试的验证思路与实现细节。测试套件概览验证什么、如何组织TieredStorageST是 Strimzi 系统测试systemtest中专门覆盖分层存储集成场景的套件对应文档 development-docs/systemtests/io.strimzi.systemtest.kafka.TieredStorageST.md测试实现位于 systemtest/src/test/java/io/strimzi/systemtest/kafka/TieredStorageST.java。该套件在源码中通过 JUnit 标签声明了执行范围与归属MicroShiftNotSupported(We are using Kaniko and OpenShift builds to build Kafka image with TS. To make it working on Microshift we will invest much time with not much additional value.) Tag(REGRESSION) Tag(TIERED_STORAGE)两个注解标签的含义Tag(REGRESSION)属于回归测试范围验证分层存储功能不会随版本演进而退化Tag(TIERED_STORAGE)可单独通过该标签过滤运行分层存储相关用例MicroShiftNotSupported由于需要借助 Kaniko 或 OpenShift Build 构建带插件的 Kafka 镜像该套件不支持 MicroShift 环境。套件归属 kafka 测试标签组与动态配置、监听器、节点池、版本升级、配额等用例一同保障 Kafka 核心功能在生产负载下的可靠性。套件包含的测试用例用例远程存储后端核心验证目标testTieredStorageWithAivenFileSystemPluginNFS本地文件系统类存储日志段上传至 NFS、本地删除、消费、远程删除testTieredStorageWithAivenS3PluginSeaweedFSS3 兼容对象存储日志段上传至 S3、本地删除、消费、远程删除两个用例共享同一套验证骨架上传 → 本地回收 → 远程读取 → 远程删除差别仅在于远程存储实现文件系统 vs S3 对象存储以及随之而来的基础设施部署NFS vs SeaweedFS。测试前置条件五个必须的准备步骤根据套件文档正式执行测试前需要依次完成以下准备工作步骤操作预期结果1创建测试命名空间命名空间创建成功2基于镜像全名、基础镜像、Dockerfile 路径等参数构建 Kafka 镜像通过 Kaniko 或 OpenShift Build并集成 Aiven Tiered Storage 插件构建出集成 Aiven Tiered Storage 插件的 Kafka 镜像3在测试命名空间部署 SeaweedFS并在 SeaweedFS Pod 内初始化 IAMSeaweedFS 部署完成且 IAM 初始化完毕4在 SeaweedFS 中初始化测试用 BucketBucket 初始化完成5部署 Cluster OperatorCluster Operator 部署完成在源码中这些前置逻辑集中于setup()方法与辅助方法中BeforeAll void setup() throws IOException { // 先安装 Cluster OperatorRBAC 资源创建在 CO 命名空间避免测试命名空间被删除 SetupClusterOperator.getInstance().withDefaultConfiguration().install(); suiteStorage new TestStorage(KubeResourceManager.get().getTestContext()); resolveTieredStorageImage(); }1. 构建带 Tiered Storage 插件的 Kafka 镜像这是整个测试最关键的前置步骤。Strimzi 的 Kafka 镜像本身不包含分层存储插件必须通过自定义镜像把插件打进/opt/kafka/plugins/目录。测试用 Dockerfile 位于 systemtest/src/test/resources/tiered-storage/DockerfileARG BASE_IMAGE FROM ${BASE_IMAGE} ARG AIVEN_PLUGIN_VERSION1.1.1 ARG TIERED_STORAGE_URLhttps://github.com/Aiven-Open/tiered-storage-for-apache-kafka/releases/download USER root:root RUN mkdir -p /opt/kafka/plugins/tiered-storage RUN curl -sL $TIERED_STORAGE_URL/v$AIVEN_PLUGIN_VERSION/s3-$AIVEN_PLUGIN_VERSION.tgz | tar -xz --strip-components1 -C /opt/kafka/plugins/tiered-storage RUN curl -sL $TIERED_STORAGE_URL/v$AIVEN_PLUGIN_VERSION/core-$AIVEN_PLUGIN_VERSION.tgz | tar -xz --strip-components1 -C /opt/kafka/plugins/tiered-storage RUN curl -sL $TIERED_STORAGE_URL/v$AIVEN_PLUGIN_VERSION/filesystem-$AIVEN_PLUGIN_VERSION.tgz | tar -xz --strip-components1 -C /opt/kafka/plugins/tiered-storage RUN rm -rf /tmp/tiered-storage USER 1001该 Dockerfile 的要点以BASE_IMAGE默认quay.io/strimzi/kafka:latest-kafka-最新支持的 Kafka 版本为基础镜像下载并解压 Aiven Tiered Storage 插件的core、s3、filesystem三个压缩包到/opt/kafka/plugins/tiered-storage最后将运行用户切回 UID 1001Kafka 容器非 root 运行。镜像构建策略在resolveTieredStorageImage()中实现支持两种方式private void resolveTieredStorageImage() throws IOException { if (Environment.KAFKA_TIERED_STORAGE_IMAGE.isEmpty()) { // 方式一用 Kaniko / OpenShift Build 现场构建 ImageBuild.buildImage(suiteStorage.getNamespaceName(), IMAGE_NAME, TIERED_STORAGE_DOCKERFILE, BUILT_IMAGE_TAG, Environment.KAFKA_TIERED_STORAGE_BASE_IMAGE); tieredStorageImageName Environment.getImageOutputRegistry(suiteStorage.getNamespaceName(), IMAGE_NAME, BUILT_IMAGE_TAG); } else { // 方式二直接使用预构建镜像 tieredStorageImageName Environment.KAFKA_TIERED_STORAGE_IMAGE; } }相关环境变量定义在 systemtest/src/main/java/io/strimzi/systemtest/Environment.java并在 development-docs/TESTING.md 中有完整说明环境变量作用默认值KAFKA_TIERED_STORAGE_IMAGE已包含 Aiven Tiered Storage 插件的 Kafka 镜像。若同时配置了KAFKA_TIERED_STORAGE_BASE_IMAGE本变量优先不再构建镜像空KAFKA_TIERED_STORAGE_CLASSPATHKafka 镜像内 Tiered Storage 插件的类路径会写入KafkaCR 的classPath字段/opt/kafka/plugins/tiered-storage/*KAFKA_TIERED_STORAGE_BASE_IMAGE用于构建新镜像的基础 Kafka 镜像quay.io/strimzi/kafka:latest-kafka-最新支持的 Kafka 版本24. 部署 SeaweedFS 并初始化 BucketS3 用例需要 S3 兼容的存储后端测试选择的是部署在集群内的 SeaweedFS。相关实现位于 systemtest/src/main/java/io/strimzi/systemtest/resources/seaweedfs/SetupSeaweedFS.javapublic static final String SEAWEEDFS seaweedfs; public static final String ADMIN_CREDS seaweedfsadminLongerThan16BytesForFIPS; public static final int SEAWEEDFS_PORT 8333; private static final String SEAWEEDFS_IMAGE mirror.gcr.io/chrislusf/seaweedfs:3.99;测试用 Bucket 名为test-bucket见TieredStorageST中的常量BUCKET_NAME在用例内通过SetupSeaweedFS.createBucket(...)创建private void deploySeaweedFSInstance() { SetupSeaweedFS.deploySeaweedFS(suiteStorage.getNamespaceName()); SetupSeaweedFS.createBucket(suiteStorage.getNamespaceName(), BUCKET_NAME); }SeaweedFS 的 S3 服务地址在 Kafka CR 中以http://seaweedfs.namespace.svc.cluster.local:8333的形式引用详见下文 S3 用例配置。用例一testTieredStorageWithAivenFileSystemPluginNFS 后端该用例使用 Aiven Tiered Storage 的 FileSystem 插件远程存储为测试期间部署的 NFS 实例。完整执行步骤如下表所示步骤操作预期结果1部署 10Gi PV 的 KafkaNodePoolKafkaNodePool 按指定配置部署成功2部署 NFS 实例含 RoleBinding、ServiceAccount、Service、StorageClass 等资源NFS 相关资源部署成功3部署 Kafka CR挂载额外 NFS 卷Tiered Storage 配置指向 NFS 路径使用构建好的 Kafka 镜像调小remote.log.manager.task.interval.ms与log.retention.check.interval.ms以加速日志上传与本地删除Kafka CR 部署成功间隔优化生效4创建开启分层存储同步的 Topicsegment 大小设为 10mb加速同步Topic 创建成功5启动持续生产者向 Kafka 发送数据生产者开始持续发送数据6等待 NFS 大小超过一个日志段大小即已收到 Kafka 数据NFS 中至少包含一个来自 Kafka 的日志段7等待 earliest-local offset 大于 0已上传到 NFS 的日志段在本地被删除8启动消费者消费全部已生产消息部分消息应位于 NFS 中消费者成功消费全部消息9修改 Topic 配置为retention.ms10s以测试远程日志删除Topic 配置修改成功10等待 NFS 数据被删除NFS 中的数据被删除NFS 实例的部署方式deployNfsInstance()方法通过kubectl apply的方式应用 systemtest/src/test/resources/nfs/nfs.yaml并做两处关键处理// 允许 NetworkPolicy 放行 NFS 流量在 default to deny all 模式下必需 NetworkPolicyUtils.allowNetworkPolicyAllIngressForMatchingLabel(suiteStorage.getNamespaceName(), nfs, Map.of(TestConstants.APP_POD_LABEL, nfs-server-provisioner)); // 将模板中的命名空间占位符替换为实际测试命名空间 String instanceYamlContent ReadWriteUtils.readFile(NFS_INSTANCE_PATH).replace(NAMESPACE_TO_BE_CHANGE, suiteStorage.getNamespaceName());NFS 实例的 Pod 以test-nfs-server-provisioner为前缀命名部署完成后等待其就绪StatefulSetUtils.waitForAllStatefulSetPodsReady(suiteStorage.getNamespaceName(), test-nfs-server-provisioner, 1);NFS 卷通过名为nfs-pvc的 PVC 暴露给 Kafka Pod常量NfsUtils.NFS_PVC_NAME见 systemtest/src/main/java/io/strimzi/systemtest/utils/specific/NfsUtils.java。部署 Kafka 集群与 FileSystem 后端分层存储配置首先部署 KafkaNodePool3 个 broker 节点使用 10Gi 持久化存储deleteClaim: true1 个 controller 节点KafkaNodePoolTemplates.brokerPoolPersistentStorage(suiteStorage.getNamespaceName(), testStorage.getBrokerPoolName(), testStorage.getClusterName(), 3) .editSpec() .withNewPersistentClaimStorage() .withSize(10Gi) .withDeleteClaim(true) .endPersistentClaimStorage() .endSpec() .build(), KafkaNodePoolTemplates.controllerPoolPersistentStorage(suiteStorage.getNamespaceName(), testStorage.getControllerPoolName(), testStorage.getClusterName(), 1).build()随后部署KafkaCR核心分层存储配置如下与文档步骤 3 对应KafkaTemplates.kafka(suiteStorage.getNamespaceName(), testStorage.getClusterName(), 3) .editSpec() .editKafka() .withImage(tieredStorageImageName) // 使用构建好的带插件镜像 .withNewTieredStorageCustomTiered() .withNewRemoteStorageManager() .withClassName(io.aiven.kafka.tieredstorage.RemoteStorageManager) .withClassPath(Environment.KAFKA_TIERED_STORAGE_CLASSPATH) // /opt/kafka/plugins/tiered-storage/* .addToConfig(storage.backend.class, io.aiven.kafka.tieredstorage.storage.filesystem.FileSystemStorage) .addToConfig(storage.root, MOUNT_PATH) // /mnt/nfs .addToConfig(chunk.size, 4194304) // 4MiB .endRemoteStorageManager() .endTieredStorageCustomTiered() // 为 Kafka Pod 额外挂载 NFS 卷 .withNewTemplate() .withNewPod() .addNewVolume() .withName(VOLUME_NAME) // nfs-volume .withNewPersistentVolumeClaim(NFS_PVC_NAME, false) // nfs-pvc .endVolume() .endPod() .withNewKafkaContainer() .addToVolumeMounts(volumeMounts) // 挂载到 /mnt/nfs .endKafkaContainer() .endTemplate() // 缩短调度间隔以加速测试 .addToConfig(remote.log.manager.task.interval.ms, 5000) .addToConfig(log.retention.check.interval.ms, 5000) .endKafka() .endSpec() .build()这段配置揭示了 FileSystem 后端的关键设计远程存储必须对 Kafka broker 可见。由于文件系统后端需要 broker 直接写文件Kafka Pod 必须把 NFS 卷nfs-pvc挂载到/mnt/nfs插件配置中的storage.root指向该挂载路径。这与 S3 后端走网络协议访问对象存储有本质区别。创建开启分层存储的 TopicTopic 配置体现了分层存储测试的核心参数组合与文档步骤 4 对应KafkaTopicTemplates.topic(suiteStorage.getNamespaceName(), testStorage.getTopicName(), testStorage.getClusterName()) .editSpec() .addToConfig(file.delete.delay.ms, 1000) .addToConfig(local.retention.ms, 1000) // 允许分层存储同步 .addToConfig(remote.storage.enable, true) // 字节保留设为 1024mb .addToConfig(retention.bytes, 1073741824) .addToConfig(retention.ms, 86400000) // segment 大小设为 10mb加速数据同步 .addToConfig(segment.bytes, SEGMENT_BYTE) .endSpec() .build()关键参数解析参数值作用remote.storage.enabletrue开启该 Topic 的分层存储同步是分层存储生效的开关segment.bytes10485761MiB日志段大小。测试注释说明Segment size is set to 10mb to make it quicker to sync data实际取值为 1MiB 常量SEGMENT_BYTE 1048576段越小、越早封段上传越频繁、同步验证越快local.retention.ms1000本地保留时间段上传远程后 1 秒即从本地删除加速本地删除验证file.delete.delay.ms1000文件删除延迟缩短本地段清理的等待时间retention.bytes1073741824字节保留上限确保测试期间段不会被字节级保留策略提前清除retention.ms86400000保留 24 小时测试阶段不触发删除保证先验证读取、再验证删除的顺序生产者、NFS 数据上浮验证与远程消费测试使用内部客户端 Job 形态的生产者/消费者KafkaProducerConsumer向 Kafka 发送 10,000 条、每条 300 个#字符的消息.withMessageCount(MESSAGE_COUNT) // 10_000 .withDelayMs(1) .withMessage(String.join(, Collections.nCopies(300, #)))生产者持续写入后测试等待 NFS 中的数据量超过一个日志段大小// wait for logs uploaded to NFS NfsUtils.waitForSizeInNfs(testStorage.getNamespaceName(), size - size SEGMENT_BYTE);NfsUtils的实现逻辑是在 NFS 提供者 Pod 内执行du -sb统计/export/volumeName目录大小NfsUtils.java 中getSizeOfDirectoryInPod使用du -sb path-b以字节为单位、-s只汇总总量从而精确判断远程存储是否收到了 Kafka 日志段。earliest-local offset验证本地日志已被回收本地段已删除的验证依靠 AdminClient 查询EARLIEST_LOCAL_TIMESTAMP偏移实现这是 Kafka 分层存储 API 提供的专门机制private void waitForEarliestLocalOffsetGreaterThanZero(String namespace, String adminName, String topicName) { final AdminClient adminClient AdminClientUtils.getConfiguredAdminClient(namespace, adminName); TestUtils.waitFor(earliest-local offset to be higher than 0, TestConstants.GLOBAL_POLL_INTERVAL_5_SECS, TestConstants.GLOBAL_TIMEOUT_LONG, () - { // 获取 earliest-local offsets // 当本地数据被删除时earliest-local offset 应大于 0 String offsetData adminClient.fetchOffsets(topicName, String.valueOf(ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP)); long earliestLocalOffset 0; try { earliestLocalOffset AdminClientUtils.getPartitionsOffset(offsetData, 0); LOGGER.info(earliest-local offset for topic {} is {}, topicName, earliestLocalOffset); } catch (JsonProcessingException e) { return false; } return earliestLocalOffset 0; }); }验证原理EARLIEST_LOCAL_TIMESTAMP返回的是仅存在于本地存储的最早偏移。若日志段全部上传到 NFS 且本地已删除则该值从 0 变为大于 0因为偏移 0 对应的段已在本地被清理本地最早的段起始偏移前移。源码注释明确写道Check that data are not present locally, earliest-local offset should be higher than 0。远程数据读取验证本地删除验证通过后测试启动消费者 JobKubeResourceManager.get().createResourceWithWait(kafkaProducerConsumer.getConsumer().getJob()); ClientUtils.waitForClientSuccess(testStorage.getNamespaceName(), testStorage.getConsumerName(), MESSAGE_COUNT);源码中的注释解释了这一验证的推理链// Verify we can consume messages from (a) remote storage and (b) local storage. // Because we have verified earlier that the log segments are moved to remote storage // (by SeaweedFS size check) and deleted locally (by earliest-local offset check), // we can verify (a) and (b) by checking if we can consume all messages successfully.即既然已验证段已上传远程NFS/SeaweedFS 中有数据且本地已删除earliest-local offset 0那么消费者能完整消费全部 10,000 条消息就同时证明了从远程存储读取和从本地存储读取两条路径都正常——这正是分层存储读取链路的核心验证。远程数据删除验证最后将 Topic 的retention.ms改为 1000010 秒触发远程日志删除// Delete data KafkaTopicUtils.replace( testStorage.getNamespaceName(), testStorage.getTopicName(), topic - topic.getSpec().getConfig().put(retention.ms, 10000) ); // wait for remote data deletion NfsUtils.waitForSizeInNfs(testStorage.getNamespaceName(), size - size SEGMENT_BYTE);由于 broker 侧已配置log.retention.check.interval.ms5000保留策略检查每 5 秒执行一次Topic 保留期缩短到 10 秒后远程段被清除NFS 目录大小回落至一个段大小以下。用例二testTieredStorageWithAivenS3PluginSeaweedFS 后端该用例使用 Aiven Tiered Storage 的 S3 插件远程存储为 SeaweedFS 提供的 S3 兼容端点验证步骤与 NFS 用例高度对称步骤操作预期结果1部署 10Gi PV 的 KafkaNodePoolKafkaNodePool 部署成功2部署 Kafka CRTiered Storage 配置指向 SeaweedFS S3使用构建好的 Kafka 镜像调小remote.log.manager.task.interval.ms与log.retention.check.interval.msKafka CR 部署成功间隔优化生效3创建开启分层存储同步的 Topicsegment 大小为 10mbTopic 创建成功4启动持续生产者发送数据生产者开始发送数据5等待 SeaweedFS 非空已收到 Kafka 数据SeaweedFS 包含 Kafka 数据6等待 earliest-local offset 大于 0已上传的日志段在本地被删除7启动消费者消费全部消息部分消息应位于 SeaweedFS消费者成功消费全部消息8修改 Topic 配置为retention.ms10s配置修改成功9等待 SeaweedFS 大小为 0SeaweedFS 中的数据被删除S3 后端的 Kafka CR 配置S3 用例无需挂载额外卷对象存储走网络协议配置核心如下KafkaTemplates.kafka(suiteStorage.getNamespaceName(), testStorage.getClusterName(), 3) .editSpec() .editKafka() .withImage(tieredStorageImageName) .withNewTieredStorageCustomTiered() .withNewRemoteStorageManager() .withClassName(io.aiven.kafka.tieredstorage.RemoteStorageManager) .withClassPath(Environment.KAFKA_TIERED_STORAGE_CLASSPATH) .addToConfig(storage.backend.class, io.aiven.kafka.tieredstorage.storage.s3.S3Storage) .addToConfig(chunk.size, 4194304) // s3 config .addToConfig(storage.s3.endpoint.url, http:// SetupSeaweedFS.SEAWEEDFS . suiteStorage.getNamespaceName() .svc.cluster.local: SetupSeaweedFS.SEAWEEDFS_PORT) .addToConfig(storage.s3.bucket.name, BUCKET_NAME) // test-bucket .addToConfig(storage.s3.region, us-east-1) .addToConfig(storage.s3.path.style.access.enabled, true) .addToConfig(storage.aws.access.key.id, SetupSeaweedFS.ADMIN_CREDS) .addToConfig(storage.aws.secret.access.key, SetupSeaweedFS.ADMIN_CREDS) .endRemoteStorageManager() .endTieredStorageCustomTiered() // 缩短调度间隔以加速测试 .addToConfig(remote.log.manager.task.interval.ms, 5000) .addToConfig(log.retention.check.interval.ms, 5000) .endKafka() .endSpec() .build()S3 插件配置项与生产环境对照配置项测试值生产对应说明storage.backend.classio.aiven.kafka.tieredstorage.storage.s3.S3StorageS3 存储后端实现类storage.s3.endpoint.urlhttp://seaweedfs.ns.svc.cluster.local:8333S3 端点生产环境指向真实 S3/兼容对象存储SeaweedFS 端口固定为 8333storage.s3.bucket.nametest-bucket存储 Bucketstorage.s3.regionus-east-1S3 区域storage.s3.path.style.access.enabledtrue使用 path-style 访问SeaweedFS 等兼容存储通常要求开启storage.aws.access.key.id/storage.aws.secret.access.keyseaweedfsadminLongerThan16BytesForFIPS访问凭据测试中读写一致chunk.size41943044MiB对象分块大小决定上传到对象存储的单个对象尺寸SeaweedFS 数据上浮与删除验证上传验证通过SeaweedFSUtils.waitForDataInSeaweedFS实现其核心是查询 Bucket 内对象数量见 SeaweedFSUtils.javapublic static void waitForDataInSeaweedFS(String namespaceName, String bucketName) { TestUtils.waitFor(data sync from Kafka to SeaweedFS, TestConstants.GLOBAL_POLL_INTERVAL_MEDIUM, TestConstants.GLOBAL_TIMEOUT_LONG, () - { long objectCount getBucketSize(namespaceName, bucketName); LOGGER.info(Bucket {} contains {} objects, bucketName, objectCount); return objectCount 0; }); }删除验证则等待对象数归零public static void waitForNoDataInSeaweedFS(String namespaceName, String bucketName) { TestUtils.waitFor(data deletion in SeaweedFS, TestConstants.GLOBAL_POLL_INTERVAL_MEDIUM, TestConstants.GLOBAL_TIMEOUT_LONG, () - { long objectCount getBucketSize(namespaceName, bucketName); LOGGER.info(Bucket {} contains {} chunks, bucketName, objectCount); return objectCount 0; }); }其余环节earliest-local offset 验证、消费验证、retention.ms修改触发远程删除与 NFS 用例完全一致共同覆盖 S3 对象存储路径的完整生命周期。从源码看 Strimzi 的分层存储 API 模型测试中使用的tieredStorage配置对应KafkaCR 的 API 模型源码位于 api/src/main/java/io/strimzi/api/kafka/model/kafka/tieredstorage/。TieredStorage抽象基类与类型判别TieredStorage.java 是抽象基类通过 Jackson 的JsonTypeInfo按type属性进行多态反序列化JsonTypeInfo( use JsonTypeInfo.Id.NAME, include JsonTypeInfo.As.EXISTING_PROPERTY, property type ) JsonSubTypes({ JsonSubTypes.Type(value TieredStorageCustom.class, name TieredStorage.TYPE_CUSTOM)} ) public abstract class TieredStorage implements UnknownPropertyPreserving { public static final String TYPE_CUSTOM custom; Description(Storage type, only custom is supported at the moment.) public abstract String getType(); ... }从源码可以看出当前 Strimzi 的tieredStorage.type只支持custom一种取值。TieredStorageCustom 与 RemoteStorageManagerTieredStorageCustom.java 是custom类型的实现包含一个remoteStorageManager字段。RemoteStorageManager.java 定义了三个核心字段其 Javadoc 直接说明了与 Kafka broker 配置的映射关系字段说明classNameRemoteStorageManager实现类的全限定名classPath插件实现类的类路径测试中为/opt/kafka/plugins/tiered-storage/*config传给RemoteStorageManager实现的额外配置键会自动加rsm.config.前缀并追加到 Kafka broker 配置config字段的映射逻辑尤为关键源码注释Description(The additional configuration map for the RemoteStorageManager implementation. Keys will be automatically prefixed with rsm.config., and added to Kafka broker configuration.)也就是说测试中写入的storage.s3.bucket.name、storage.backend.class等键最终会以rsm.config.storage.s3.bucket.name等形式注入 broker 配置由 Aiven 插件的RemoteStorageManager读取。Strimzi 官方文档中的分层存储配置指南仓库文档 documentation/modules/configuring/ref-storage-tiered.adoc 提供了与测试配置一致、面向生产用户的配置模板可作为理解测试配置与生产配置映射关系的参考apiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name: my-cluster spec: kafka: tieredStorage: type: custom # 1) 类型必须为 custom remoteStorageManager: # 2) 自定义 RemoteStorageManager 实现配置 className: com.example.kafka.tiered.storage.s3.S3RemoteStorageManager classPath: /opt/kafka/plugins/tiered-storage-s3/* config: storage.bucket.name: my-bucket # 3) 自动加 rsm.config. 前缀传给实现 config: rlmm.config.remote.log.metadata.topic.replication.factor: 1 # 4) RLMM 专属配置需加 rlmm.config. 前缀该文档补充了几个测试源码未直接体现的关键事实RLMMRemote Log Metadata ManagerStrimzi 开启自定义分层存储时使用 Kafka 的TopicBasedRemoteLogMetadataManager管理远程存储元数据RLMM 专属配置必须加rlmm.config.前缀写入spec.kafka.config例如rlmm.config.remote.log.metadata.topic.replication.factor插件前置条件使用自定义分层存储前必须通过构建自定义容器镜像把插件打入 Strimzi Kafka 镜像生产就绪状态分层存储是 Kafka 3.9.0 起的生产级特性Strimzi 同样支持引入前应审查 Apache Kafka 官方文档列出的分层存储限制。存储选型边界documentation/modules/configuring/con-considerations-for-data-storage.adoc 中明确了分层存储在整体存储架构中的定位主 broker 存储持久卷、JBOD承载近期数据远程分层存储如对象存储承载历史数据Strimzi 以块存储为主存储的推荐类型如 AWS EBS、Azure Disk、GCP Persistent Disk文件系统型存储如 NFS不保证作为主存储的稳定性分层存储是附加能力配置好主存储后可在集群级配置分层存储并通过 Topic 级remote.storage.enable对特定 Topic 启用自定义插件投入生产前必须验证其性能与兼容性要求。相关文档入口documentation/assemblies/configuring/assembly-storage.adoc 将分层存储作为Configuring Kafka storage章节的组成部分与临时存储、持久存储、JBOD 并列阐述API 参考见 documentation/api/io.strimzi.api.kafka.model.kafka.tieredstorage.TieredStorageCustom.adoc。测试调优参数总结两个用例都通过缩短 broker 侧两个调度间隔来加速测试这是理解分层存储内部时序的关键参数测试值默认意义影响remote.log.manager.task.interval.ms5000远程日志管理器的任务调度间隔控制日志段上传到远程存储的频率调小加速上传验证log.retention.check.interval.ms5000日志保留检查间隔控制本地/远程日志删除检查频率调小加速删除验证生产环境无需如此激进可按实际负载调整。测试通过压缩时序换取验证速度同时完整覆盖了分层存储的四个核心能力上传Upload日志段从本地同步到远程存储NFS/SeaweedFS 出现数据本地回收Local deletion上传完成后本地段被删除earliest-local offset 0远程读取Remote read消费者能从远程存储取回全部消息消费 10,000 条成功远程删除Remote deletion保留策略到期后远程数据被清除NFS/SeaweedFS 数据归零。如何运行这套测试测试位于 systemtest 模块运行前需满足 development-docs/TESTING.md 中的系统测试环境要求可用 Kubernetes 集群、Cluster Operator 部署权限、镜像构建能力等。针对分层存储套件可按 JUnit 标签过滤运行# 仅运行分层存储相关用例 mvn -f systemtest/pom.xml verify -Dgroupstiered-storage或直接针对测试类运行mvn -f systemtest/pom.xml verify -Dit.testTieredStorageST运行前根据环境设置镜像相关变量不设置则自动执行镜像构建# 可选直接指定已含插件的镜像跳过构建 export KAFKA_TIERED_STORAGE_IMAGEregistry.example.com/strimzi/kafka-tiered-storage:latest # 可选指定构建基础镜像KAFKA_TIERED_STORAGE_IMAGE 未设置时生效 export KAFKA_TIERED_STORAGE_BASE_IMAGEquay.io/strimzi/kafka:latest-kafka-3.9.0 # 可选插件类路径默认 /opt/kafka/plugins/tiered-storage/* export KAFKA_TIERED_STORAGE_CLASSPATH/opt/kafka/plugins/tiered-storage/*小结TieredStorageST以Aiven 插件 两种远程存储后端的组合为 Strimzi 分层存储功能提供了端到端的回归保障FileSystem 用例验证了broker 本地挂载远程卷的文件系统路径S3 用例验证了网络访问对象存储的通用路径。测试背后的三个关键机制值得在生产部署中复用CR 配置tieredStorage.type: customremoteStorageManagerclassName/classPath/configconfig 键自动加rsm.config.前缀注入 brokerTopic 开关remote.storage.enable: true配合segment.bytes、local.retention.ms等参数控制同步行为验证手段EARLIEST_LOCAL_TIMESTAMP偏移检查本地回收状态、远程存储容量/对象数检查上传与删除状态二者结合形成完整的生命周期闭环。如果你正在规划 Kafka 集群的分层存储落地这套测试的配置模板含调优参数可以直接作为KafkaCR 与 Topic 配置的起点再结合 ref-storage-tiered.adoc 中的生产配置指南与 RLMM 说明进行调整。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表