You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何仅用主题部分字段创建GlobalKTable以节省存储?

问题:Kafka Streams GlobalKTable 状态存储空间优化

现有实现与痛点

我们用Kafka Streams处理输入流和存储客户端数据的压缩主题,通过GlobalKTable关联二者,代码如下:

StreamsBuilder builder = new StreamsBuilder();

GlobalKTable<String, Client> clients
    = builder.globalTable(CLIENT_TOPIC, Consumed.with(Serdes.String(), clientSerde));
        
KStream<String, Foo> foos = builder.stream(FOO_TOPIC);

KStream<String, Bar> bars = foos
                .leftJoin(
                        clients,
                        (streamKey, streamValue) -> streamValue.getClientId().toString(),
                        new FooClientJoiner()
                );

这段代码能正常运行,但整个CLIENT_TOPIC的数据会被存入RocksDB状态存储(硬盘),占用大量空间,已经成为问题。实际上我们只需要客户端数据中的某一个字段,原本期望能通过流转换后生成只包含所需字段的GlobalKTable,比如:

StreamsBuilder builder = new StreamsBuilder();

// 期望从流生成精简的GlobalKTable,但目前无法实现
GlobalKTable<String, ReducedClient> reducedClients
    = builder.stream(CLIENT_TOPIC)
        .map((key,value) -> new KeyValue<String, String>(key, value.getTheOneColumnINeed()))
        .toGlobalKTable();
        
KStream<String, Foo> foos = builder.stream(FOO_TOPIC);

KStream<String, Bar> bars = foos
                .leftJoin(
                        reducedClients,
                        (streamKey, streamValue) -> streamValue.getClientId().toString(),
                        new FooReducedClientsJoiner()
                );

但目前Kafka Streams不支持直接从流创建GlobalKTable,请问有没有等效方案?或者在仅需主题小部分数据的情况下,如何减少CLIENT_TOPIC缓存的空间占用?

已尝试的无效方案:

  • 手动修改生成的Client类,仅保留所需字段,但GlobalKTable存储大小无变化
  • 调整RocksDB的压缩设置,无明显效果

可行解决方案

方案1:自定义SerDe,仅序列化/反序列化所需字段

不需要修改Client类,而是实现一个自定义SerDe,在反序列化CLIENT_TOPIC消息时,只提取需要的单个字段,序列化时也只处理该字段。这样GlobalKTable存储的就只有这个字段的数据,而非整个Client对象。

示例代码:

// 自定义SerDe,仅处理需要的字段
public class ReducedClientSerde extends Serdes.WrapperSerde<String> {
    public ReducedClientSerde() {
        super(new StringSerializer(), new Deserializer<String>() {
            @Override
            public String deserialize(String topic, byte[] data) {
                // 反序列化时,只提取Client对象中的目标字段
                Client client = originalClientSerde.deserializer().deserialize(topic, data);
                return client.getTheOneColumnINeed();
            }
        });
    }
}

// 使用自定义SerDe创建GlobalKTable
GlobalKTable<String, String> reducedClients = builder.globalTable(
    CLIENT_TOPIC,
    Consumed.with(Serdes.String(), new ReducedClientSerde())
);

// 后续join逻辑调整为使用单个字段
KStream<String, Bar> bars = foos
    .leftJoin(
        reducedClients,
        (streamKey, streamValue) -> streamValue.getClientId().toString(),
        (foo, clientField) -> new Bar(foo, clientField) // 调整Joiner逻辑
    );

方案2:提前处理CLIENT_TOPIC,生成精简主题

在Kafka Streams应用之外,或者单独开一个Streams子任务,先消费CLIENT_TOPIC,提取所需字段后写入一个新的精简主题(比如CLIENT_REDUCED_TOPIC),然后在主逻辑中用这个精简主题创建GlobalKTable。

示例代码(精简主题生成逻辑):

// 单独的流处理逻辑,生成精简主题
StreamsBuilder builder = new StreamsBuilder();
builder.stream(CLIENT_TOPIC)
    .map((key, client) -> new KeyValue<>(key, client.getTheOneColumnINeed()))
    .to("CLIENT_REDUCED_TOPIC", Produced.with(Serdes.String(), Serdes.String()));

// 主逻辑中使用精简主题创建GlobalKTable
GlobalKTable<String, String> reducedClients = builder.globalTable(
    "CLIENT_REDUCED_TOPIC",
    Consumed.with(Serdes.String(), Serdes.String())
);

这种方案的好处是主逻辑更简洁,且精简主题可以被多个应用复用,缺点是多了一个主题和额外的处理任务。

方案3:使用Transformer替代GlobalKTable(不推荐,仅特殊场景)

如果客户端数据更新不频繁,可以自定义Transformer,在其中维护一个本地缓存(比如Guava Cache),只缓存所需字段。但这种方案需要自己处理缓存过期、数据更新同步等问题,复杂度较高,不如前两种方案可靠。


关键说明

  • 自定义SerDe是最直接的方案,不需要额外主题,直接在GlobalKTable创建时精简数据
  • 若多个应用都需要该精简字段,提前生成精简主题的扩展性更好
  • 之前修改Client类无效,可能是因为原SerDe仍在序列化整个对象(比如用了自动生成的SerDe,比如Avro的话Schema没更新),自定义SerDe才能从根本上减少存储的数据量

内容的提问来源于stack exchange,提问作者user992990

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 01:55:16