Kafka Streams对比同键历史值转发更新及broker配置咨询
需求说明
现有两个Kafka主题:
- Streams_Input_Topic
- Streams_Output_Topic
需要基于Java开发Kafka Streams应用,实现以下逻辑:
- 从Streams_Input_Topic消费消息
- 借助持久化状态存储,校验当前消息key是否存在历史值
- 若存在历史值,对比新旧值的字段差异,只要存在字段更新,就将携带更新后完整字段的消息以原key转发到Streams_Output_Topic
- 若不存在历史值,将当前key-value存入状态存储,不做转发
逻辑示例:
// 输入消息1(首次收到key1,存入状态存储,不转发) {"key":"key1", "value":{"prop1":"value1","prop2":"value2"}} // 输入消息2(再次收到key1,对比发现prop2从value2更新为value4,触发转发) {"key":"key1", "value":{"prop1":"value1","prop2":"value4"}}
触发转发时,写入Streams_Output_Topic的消息内容如下:
{"key":"key1", "value":{"prop1":"value1","prop2":"value4"}}
现有实现代码
当前编写的代码存在逻辑错误,且未正确配置Kafka broker连接参数,代码如下:
private void prepareStream(String topic, String appId) throws IOException { System.out.println("Inside the prepare stream"); String storeName = store + "-" + topic; Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, appId); Properties p = account.connect(); props.putAll(p); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); Path path = Paths.get("persistentKeyValueStore/kafka_streams_custom_" + topic); try { path = Files.createDirectories(path); } catch (IOException e) { e.printStackTrace(); } LOG.info("Path : " + path.toAbsolutePath().toString()); props.put(StreamsConfig.STATE_DIR_CONFIG, path.toAbsolutePath().toString()); KeyValueBytesStoreSupplier stateStore = Stores.persistentKeyValueStore(storeName); LOG.debug("Creating Stream Object"); StreamsBuilder builder = new StreamsBuilder(); StoreBuilder<KeyValueStore<String, String>> keyValueStoreBuilder = Stores.keyValueStoreBuilder(Stores.persistentKeyValueStore(storeName), Serdes.String(), Serdes.String()); builder.addStateStore(keyValueStoreBuilder); builder.stream("Streams_Input_Topic", Consumed.with(Serdes.String(), Serdes.String())) .transformValues(() -> new ValueTransformerWithKey<String, String, String>() { private KeyValueStore<String, String> state; @Override public void init(final ProcessorContext context) { KafkaStreamLog.printConsole("Inside the init method of processor"); state = (KeyValueStore<String, String>) context.getStateStore(storeName); } @Override public String transform(final String key, final String value) { String prevValue = state.get(key); KafkaStreamLog.printConsole("Prev value : " + prevValue); // 原代码日志打印错误,当前值被错写为prevValue KafkaStreamLog.printConsole("Curr value : " + prevValue); if (prevValue != null) { // 原逻辑错误:此处返回了旧值,不符合转发新值的需求 return prevValue; } else { state.put(key, value); } return null; } @Override public void close() { } }, storeName).to("Streams_Output_Topic"); }
现存问题
- 未明确Kafka broker连接参数的配置方式,无法连接指定集群
- 现有转换逻辑存在错误,无法实现「仅当值更新时转发新值」的需求
- 缺少KafkaStreams实例启动的完整流程,代码无法直接运行
可运行完整实现
核心配置说明
Kafka broker连接参数直接在Streams配置的Properties中指定,固定配置项为StreamsConfig.BOOTSTRAP_SERVERS_CONFIG,值为集群broker地址列表,格式为host1:port1,host2:port2,无需单独给KStream传参,所有配置会在构建KafkaStreams实例时全局生效。
以下是修正后的完整代码,依赖Jackson做JSON字段级对比(也可替换为FastJSON等其他JSON库):
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.Produced; import org.apache.kafka.streams.kstream.ValueTransformerWithKey; import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; import java.util.Properties; public class KafkaStreamDiffService { private static final String INPUT_TOPIC = "Streams_Input_Topic"; private static final String OUTPUT_TOPIC = "Streams_Output_Topic"; private static final String STATE_STORE_NAME = "message-diff-store"; // 替换为实际的Kafka broker地址 private static final String BOOTSTRAP_SERVERS = "your-kafka-broker:9092"; private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); public void startStream() throws IOException { Properties props = new Properties(); // 配置应用ID,同一消费组内唯一 props.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-stream-diff-app"); // 配置Kafka broker连接地址 props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS); // 配置默认key/value序列化器 props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); // 配置状态存储本地路径 Path statePath = Paths.get("persistentKeyValueStore/kafka_streams_diff"); Files.createDirectories(statePath); props.put(StreamsConfig.STATE_DIR_CONFIG, statePath.toAbsolutePath().toString()); // 可选配置:提交间隔等 props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); StreamsBuilder builder = new StreamsBuilder(); // 注册持久化状态存储 StoreBuilder<KeyValueStore<String, String>> storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore(STATE_STORE_NAME), Serdes.String(), Serdes.String() ); builder.addStateStore(storeBuilder); // 构建流处理逻辑 builder.stream(INPUT_TOPIC, Consumed.with(Serdes.String(), Serdes.String())) .transformValues(() -> new ValueTransformerWithKey<String, String, String>() { private KeyValueStore<String, String> stateStore; @Override public void init(ProcessorContext context) { stateStore = context.getStateStore(STATE_STORE_NAME); } @Override public String transform(String key, String currentValue) { if (key == null || currentValue == null) { return null; } String oldValue = stateStore.get(key); // 首次收到该key,存入存储后不转发 if (oldValue == null) { stateStore.put(key, currentValue); return null; } // 新旧值字符串完全一致,不转发 if (oldValue.equals(currentValue)) { return null; } // 对比JSON字段差异,存在更新则更新存储、返回新值触发转发 try { JsonNode oldNode = OBJECT_MAPPER.readTree(oldValue); JsonNode newNode = OBJECT_MAPPER.readTree(currentValue); if (!oldNode.equals(newNode)) { stateStore.put(key, currentValue); return currentValue; } } catch (Exception e) { // JSON解析失败,默认判定为值更新,转发新值 stateStore.put(key, currentValue); return currentValue; } // 无字段差异,不转发 return null; } @Override public void close() {} }, STATE_STORE_NAME) .to(OUTPUT_TOPIC, Produced.with(Serdes.String(), Serdes.String())); // 启动流实例 KafkaStreams streams = new KafkaStreams(builder.build(), props); // 优雅关闭钩子 Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); streams.start(); } public static void main(String[] args) throws IOException { new KafkaStreamDiffService().startStream(); } }
关键说明
- 修正了原代码的逻辑错误:首次收到key仅存值不转发;新旧值一致不转发;存在字段差异时更新状态存储、转发最新值。
- 实现中使用Jackson做JSON节点级对比,避免字符串格式差异(如字段顺序不同)导致的误判,可根据自身技术栈替换为其他JSON处理库。如果消息体为非JSON格式,可直接用
String.equals()做值对比。 - 代码增加了空值校验、异常兜底、JVM关闭时的流实例优雅停服逻辑,可直接运行。
内容的提问来源于stack exchange,提问作者Patan
相关产品推荐
相关产品推荐

