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 插件,并借助 SeaweedFS(S3 兼容对象存储)与 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 核心功能在生产负载下的可靠性。
套件包含的测试用例
| 用例 | 远程存储后端 | 核心验证目标 |
|---|---|---|
testTieredStorageWithAivenFileSystemPlugin | NFS(本地文件系统类存储) | 日志段上传至 NFS、本地删除、消费、远程删除 |
testTieredStorageWithAivenS3Plugin | SeaweedFS(S3 兼容对象存储) | 日志段上传至 S3、本地删除、消费、远程删除 |
两个用例共享同一套验证骨架:上传 → 本地回收 → 远程读取 → 远程删除,差别仅在于远程存储实现(文件系统 vs S3 对象存储)以及随之而来的基础设施部署(NFS vs SeaweedFS)。
测试前置条件:五个必须的准备步骤
根据套件文档,正式执行测试前需要依次完成以下准备工作:
| 步骤 | 操作 | 预期结果 |
|---|---|---|
| 1 | 创建测试命名空间 | 命名空间创建成功 |
| 2 | 基于镜像全名、基础镜像、Dockerfile 路径等参数构建 Kafka 镜像(通过 Kaniko 或 OpenShift Build),并集成 Aiven Tiered Storage 插件 | 构建出集成 Aiven Tiered Storage 插件的 Kafka 镜像 |
| 3 | 在测试命名空间部署 SeaweedFS,并在 SeaweedFS Pod 内初始化 IAM | SeaweedFS 部署完成且 IAM 初始化完毕 |
| 4 | 在 SeaweedFS 中初始化测试用 Bucket | Bucket 初始化完成 |
| 5 | 部署 Cluster Operator | Cluster Operator 部署完成 |
在源码中,这些前置逻辑集中于setup()方法与辅助方法中:
@BeforeAll void setup() throws IOException { // 先安装 Cluster Operator(RBAC 资源创建在 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/Dockerfile:
ARG BASE_IMAGE FROM ${BASE_IMAGE} ARG AIVEN_PLUGIN_VERSION="1.1.1" ARG TIERED_STORAGE_URL="https://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-components=1 -C "/opt/kafka/plugins/tiered-storage" RUN curl -sL "$TIERED_STORAGE_URL/v$AIVEN_PLUGIN_VERSION/core-$AIVEN_PLUGIN_VERSION.tgz" | tar -xz --strip-components=1 -C "/opt/kafka/plugins/tiered-storage" RUN curl -sL "$TIERED_STORAGE_URL/v$AIVEN_PLUGIN_VERSION/filesystem-$AIVEN_PLUGIN_VERSION.tgz" | tar -xz --strip-components=1 -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 1001(Kafka 容器非 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_CLASSPATH | Kafka 镜像内 Tiered Storage 插件的类路径,会写入KafkaCR 的classPath字段 | /opt/kafka/plugins/tiered-storage/* |
KAFKA_TIERED_STORAGE_BASE_IMAGE | 用于构建新镜像的基础 Kafka 镜像 | quay.io/strimzi/kafka:latest-kafka-<最新支持的 Kafka 版本> |
2~4. 部署 SeaweedFS 并初始化 Bucket
S3 用例需要 S3 兼容的存储后端,测试选择的是部署在集群内的 SeaweedFS。相关实现位于 systemtest/src/main/java/io/strimzi/systemtest/resources/seaweedfs/SetupSeaweedFS.java:
public 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 用例配置。
用例一:testTieredStorageWithAivenFileSystemPlugin(NFS 后端)
该用例使用 Aiven Tiered Storage 的 FileSystem 插件,远程存储为测试期间部署的 NFS 实例。完整执行步骤如下表所示:
| 步骤 | 操作 | 预期结果 |
|---|---|---|
| 1 | 部署 10Gi PV 的 KafkaNodePool | KafkaNodePool 按指定配置部署成功 |
| 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 | 创建开启分层存储同步的 Topic,segment 大小设为 10mb(加速同步) | Topic 创建成功 |
| 5 | 启动持续生产者向 Kafka 发送数据 | 生产者开始持续发送数据 |
| 6 | 等待 NFS 大小超过一个日志段大小(即已收到 Kafka 数据) | NFS 中至少包含一个来自 Kafka 的日志段 |
| 7 | 等待 earliest-local offset 大于 0 | 已上传到 NFS 的日志段在本地被删除 |
| 8 | 启动消费者消费全部已生产消息,部分消息应位于 NFS 中 | 消费者成功消费全部消息 |
| 9 | 修改 Topic 配置为retention.ms=10s以测试远程日志删除 | 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 后端分层存储配置
首先部署 KafkaNodePool:3 个 broker 节点使用 10Gi 持久化存储(deleteClaim: true),1 个 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 后端(走网络协议访问对象存储)有本质区别。
创建开启分层存储的 Topic
Topic 配置体现了分层存储测试的核心参数组合(与文档步骤 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.enable | true | 开启该 Topic 的分层存储同步,是分层存储生效的开关 |
segment.bytes | 1048576(1MiB) | 日志段大小。测试注释说明"Segment size is set to 10mb to make it quicker to sync data"(实际取值为 1MiB 常量SEGMENT_BYTE = 1048576),段越小、越早封段,上传越频繁、同步验证越快 |
local.retention.ms | 1000 | 本地保留时间,段上传远程后 1 秒即从本地删除,加速"本地删除"验证 |
file.delete.delay.ms | 1000 | 文件删除延迟,缩短本地段清理的等待时间 |
retention.bytes | 1073741824 | 字节保留上限,确保测试期间段不会被字节级保留策略提前清除 |
retention.ms | 86400000 | 保留 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"。
远程数据读取验证
本地删除验证通过后,测试启动消费者 Job:
KubeResourceManager.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改为 10000(10 秒),触发远程日志删除:
// 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.ms=5000,保留策略检查每 5 秒执行一次;Topic 保留期缩短到 10 秒后,远程段被清除,NFS 目录大小回落至一个段大小以下。
用例二:testTieredStorageWithAivenS3Plugin(SeaweedFS 后端)
该用例使用 Aiven Tiered Storage 的 S3 插件,远程存储为 SeaweedFS 提供的 S3 兼容端点,验证步骤与 NFS 用例高度对称:
| 步骤 | 操作 | 预期结果 |
|---|---|---|
| 1 | 部署 10Gi PV 的 KafkaNodePool | KafkaNodePool 部署成功 |
| 2 | 部署 Kafka CR:Tiered Storage 配置指向 SeaweedFS S3,使用构建好的 Kafka 镜像;调小remote.log.manager.task.interval.ms与log.retention.check.interval.ms | Kafka CR 部署成功,间隔优化生效 |
| 3 | 创建开启分层存储同步的 Topic,segment 大小为 10mb | Topic 创建成功 |
| 4 | 启动持续生产者发送数据 | 生产者开始发送数据 |
| 5 | 等待 SeaweedFS 非空(已收到 Kafka 数据) | SeaweedFS 包含 Kafka 数据 |
| 6 | 等待 earliest-local offset 大于 0 | 已上传的日志段在本地被删除 |
| 7 | 启动消费者消费全部消息,部分消息应位于 SeaweedFS | 消费者成功消费全部消息 |
| 8 | 修改 Topic 配置为retention.ms=10s | 配置修改成功 |
| 9 | 等待 SeaweedFS 大小为 0 | SeaweedFS 中的数据被删除 |
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.class | io.aiven.kafka.tieredstorage.storage.s3.S3Storage | S3 存储后端实现类 |
storage.s3.endpoint.url | http://seaweedfs.<ns>.svc.cluster.local:8333 | S3 端点;生产环境指向真实 S3/兼容对象存储,SeaweedFS 端口固定为 8333 |
storage.s3.bucket.name | test-bucket | 存储 Bucket |
storage.s3.region | us-east-1 | S3 区域 |
storage.s3.path.style.access.enabled | true | 使用 path-style 访问(SeaweedFS 等兼容存储通常要求开启) |
storage.aws.access.key.id/storage.aws.secret.access.key | seaweedfsadminLongerThan16BytesForFIPS | 访问凭据,测试中读写一致 |
chunk.size | 4194304(4MiB) | 对象分块大小,决定上传到对象存储的单个对象尺寸 |
SeaweedFS 数据上浮与删除验证
上传验证通过SeaweedFSUtils.waitForDataInSeaweedFS实现,其核心是查询 Bucket 内对象数量(见 SeaweedFSUtils.java):
public 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 与 RemoteStorageManager
TieredStorageCustom.java 是custom类型的实现,包含一个remoteStorageManager字段。
RemoteStorageManager.java 定义了三个核心字段,其 Javadoc 直接说明了与 Kafka broker 配置的映射关系:
| 字段 | 说明 |
|---|---|
className | RemoteStorageManager实现类的全限定名 |
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. 前缀该文档补充了几个测试源码未直接体现的关键事实:
- RLMM(Remote Log Metadata Manager):Strimzi 开启自定义分层存储时,使用 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.ms | 5000 | 远程日志管理器的任务调度间隔 | 控制日志段上传到远程存储的频率;调小加速上传验证 |
log.retention.check.interval.ms | 5000 | 日志保留检查间隔 | 控制本地/远程日志删除检查频率;调小加速删除验证 |
生产环境无需如此激进,可按实际负载调整。测试通过"压缩时序"换取验证速度,同时完整覆盖了分层存储的四个核心能力:
- 上传(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 -Dgroups=tiered-storage或直接针对测试类运行:
mvn -f systemtest/pom.xml verify -Dit.test=TieredStorageST运行前根据环境设置镜像相关变量(不设置则自动执行镜像构建):
# 可选:直接指定已含插件的镜像,跳过构建 export KAFKA_TIERED_STORAGE_IMAGE=registry.example.com/strimzi/kafka-tiered-storage:latest # 可选:指定构建基础镜像(KAFKA_TIERED_STORAGE_IMAGE 未设置时生效) export KAFKA_TIERED_STORAGE_BASE_IMAGE=quay.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: custom+remoteStorageManager(className/classPath/config),config 键自动加rsm.config.前缀注入 broker; - Topic 开关:
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),仅供参考