news 2026/9/17 15:20:34

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 插件,并借助 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 核心功能在生产负载下的可靠性。

套件包含的测试用例

用例远程存储后端核心验证目标
testTieredStorageWithAivenFileSystemPluginNFS(本地文件系统类存储)日志段上传至 NFS、本地删除、消费、远程删除
testTieredStorageWithAivenS3PluginSeaweedFS(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 内初始化 IAMSeaweedFS 部署完成且 IAM 初始化完毕
4在 SeaweedFS 中初始化测试用 BucketBucket 初始化完成
5部署 Cluster OperatorCluster 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 插件的cores3filesystem三个压缩包到/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_CLASSPATHKafka 镜像内 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 的 KafkaNodePoolKafkaNodePool 按指定配置部署成功
2部署 NFS 实例(含 RoleBinding、ServiceAccount、Service、StorageClass 等资源)NFS 相关资源部署成功
3部署 Kafka CR:挂载额外 NFS 卷,Tiered Storage 配置指向 NFS 路径,使用构建好的 Kafka 镜像;调小remote.log.manager.task.interval.mslog.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.enabletrue开启该 Topic 的分层存储同步,是分层存储生效的开关
segment.bytes1048576(1MiB)日志段大小。测试注释说明"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"。

远程数据读取验证

本地删除验证通过后,测试启动消费者 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 的 KafkaNodePoolKafkaNodePool 部署成功
2部署 Kafka CR:Tiered Storage 配置指向 SeaweedFS S3,使用构建好的 Kafka 镜像;调小remote.log.manager.task.interval.mslog.retention.check.interval.msKafka CR 部署成功,间隔优化生效
3创建开启分层存储同步的 Topic,segment 大小为 10mbTopic 创建成功
4启动持续生产者发送数据生产者开始发送数据
5等待 SeaweedFS 非空(已收到 Kafka 数据)SeaweedFS 包含 Kafka 数据
6等待 earliest-local offset 大于 0已上传的日志段在本地被删除
7启动消费者消费全部消息,部分消息应位于 SeaweedFS消费者成功消费全部消息
8修改 Topic 配置为retention.ms=10s配置修改成功
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 端口固定为 8333
storage.s3.bucket.nametest-bucket存储 Bucket
storage.s3.regionus-east-1S3 区域
storage.s3.path.style.access.enabledtrue使用 path-style 访问(SeaweedFS 等兼容存储通常要求开启)
storage.aws.access.key.id/storage.aws.secret.access.keyseaweedfsadminLongerThan16BytesForFIPS访问凭据,测试中读写一致
chunk.size4194304(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 的@JsonTypeInfotype属性进行多态反序列化:

@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 配置的映射关系:

字段说明
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.namestorage.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.ms5000远程日志管理器的任务调度间隔控制日志段上传到远程存储的频率;调小加速上传验证
log.retention.check.interval.ms5000日志保留检查间隔控制本地/远程日志删除检查频率;调小加速删除验证

生产环境无需如此激进,可按实际负载调整。测试通过"压缩时序"换取验证速度,同时完整覆盖了分层存储的四个核心能力:

  1. 上传(Upload):日志段从本地同步到远程存储(NFS/SeaweedFS 出现数据);
  2. 本地回收(Local deletion):上传完成后本地段被删除(earliest-local offset > 0);
  3. 远程读取(Remote read):消费者能从远程存储取回全部消息(消费 10,000 条成功);
  4. 远程删除(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+remoteStorageManagerclassName/classPath/config),config 键自动加rsm.config.前缀注入 broker;
  • Topic 开关remote.storage.enable: true配合segment.byteslocal.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 15:17:48

Stata面板数据实战:xtset清洗、xtreg模型与动态GMM

简介&#xff1a;面板数据是同时包含截面与时间维度的追踪数据&#xff0c;核心在于分离个体不随时间变化的异质性与时间冲击。Stata中常用xtset声明面板结构&#xff0c;再用固定效应、随机效应或组间模型进行xtreg回归&#xff0c;并通过Hausman检验判断模型取舍&#xff1b;…

作者头像 李华
网站建设 2026/9/17 15:16:31

基于SpringBoot+Vue的公交运营管理系统:三角色权限与Token鉴权

简介&#xff1a;这是一份面向高校计算机专业毕业设计场景的完整论文文档&#xff0c;主题为基于SpringBoot与Vue的城市公交运营管理系统&#xff0c;适合正在准备毕设、课程设计或需要参考前后端分离项目写法的学生与开发者。压缩包内共1个docx文件&#xff0c;约4.68MB&#…

作者头像 李华
网站建设 2026/9/17 15:15:23

ASP物流管理系统毕业设计:三层架构、数据库与IIS部署全解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华