如何仅用主题部分字段创建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

