如何利用Logstash实现Kafka Topic消息去重并转发至新Topic?
Hey, great question! Deduplicating Kafka messages based on the id field and forwarding only unique entries to Topic2 is a common requirement, and there are several robust solutions depending on your existing tech stack. Here are the most practical approaches:
1. Kafka Streams (Native Kafka Ecosystem Approach)
If you want to stay within the Kafka ecosystem, Kafka Streams is the most lightweight and integrated option. It uses persistent state stores to track processed ids, with built-in fault tolerance (state is backed up to internal Kafka topics for recovery).
Example Java Code
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.state.KeyValueStore; import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.Properties; public class KafkaDeduplicator { public static void main(String[] args) { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "deduplication-service"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); // Persistent state store to track processed IDs StoreBuilder<KeyValueStore<String, String>> idStore = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore("processed-ids"), Serdes.String(), Serdes.String()); StreamsBuilder builder = new StreamsBuilder(); builder.addStateStore(idStore); builder.stream("Topic1", Consumed.with(Serdes.String(), Serdes.String())) .filter((key, jsonValue) -> { try { // Extract ID from JSON using Jackson for reliable parsing ObjectMapper mapper = new ObjectMapper(); JsonNode node = mapper.readTree(jsonValue); String id = node.get("id").asText(); // Check if ID has been processed before KeyValueStore<String, String> store = KafkaStreams.getLocalStore("processed-ids"); if (store.get(id) == null) { store.put(id, "processed"); return true; // Keep unique message } return false; // Drop duplicate } catch (Exception e) { // Handle invalid JSON (optional: log and drop or route to dead-letter topic) return false; } }) .to("Topic2", Produced.with(Serdes.String(), Serdes.String())); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); // Graceful shutdown on app exit Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); } }
Key Notes
- Add a TTL to the state store if you only need to deduplicate recent messages (e.g., 7 days) to avoid infinite state growth.
- Kafka Streams automatically recovers state if the application restarts, so no data is lost.
2. Apache Flink (For Complex Stream Processing)
If you're already using Flink for stream processing or need advanced features like windowed deduplication, Flink is an excellent choice. It uses managed state with checkpointing for fault tolerance and scales well for large data volumes.
Example Java Code
import org.apache.flink.api.common.state.ValueState; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.KeyedProcessFunction; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer; import org.apache.flink.util.Collector; import org.apache.flink.api.java.tuple.Tuple2; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.api.common.serialization.SimpleStringSchema; import java.util.Properties; public class FlinkDeduplicator { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); // Enable checkpointing for fault tolerance Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "your-kafka-broker:9092"); kafkaProps.setProperty("group.id", "flink-deduplication-group"); // Read from Topic1 DataStream<String> inputStream = env.addSource( new FlinkKafkaConsumer<>("Topic1", new SimpleStringSchema(), kafkaProps)); // Deduplicate based on ID DataStream<String> deduplicatedStream = inputStream .map(json -> { ObjectMapper mapper = new ObjectMapper(); JsonNode node = mapper.readTree(json); return new Tuple2<>(node.get("id").asText(), json); }) .keyBy(tuple -> tuple.f0) .process(new KeyedProcessFunction<String, Tuple2<String, String>, String>() { private ValueState<Boolean> hasProcessed; @Override public void open(Configuration config) { hasProcessed = getRuntimeContext().getState( new ValueStateDescriptor<>("processed-flag", Boolean.class)); } @Override public void processElement(Tuple2<String, String> value, Context ctx, Collector<String> out) throws Exception { if (hasProcessed.value() == null) { out.collect(value.f1); hasProcessed.update(true); // Optional: Expire state after 1 day to clean up old IDs ctx.timerService().registerProcessingTimeTimer( ctx.timerService().currentProcessingTime() + 86400000); } } // Clean up expired state to prevent memory bloat @Override public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) { hasProcessed.clear(); } }); // Write unique messages to Topic2 deduplicatedStream.addSink( new FlinkKafkaProducer<>("Topic2", new SimpleStringSchema(), kafkaProps)); env.execute("Kafka Deduplication Job"); } }
Key Notes
- Flink's checkpointing ensures state is restored if the job fails mid-processing.
- You can easily adjust this logic to support windowed deduplication (e.g., only deduplicate messages within a 1-hour window) if needed.
3. Logstash + Redis (Extend Your Existing Pipeline)
Since you're already using Logstash to feed data from Oracle to Kafka, you can extend it to handle deduplication. This works best for small-scale deployments; for distributed Logstash instances, use Redis as a shared state store to avoid missing duplicates across instances.
Example Logstash Config
input { kafka { bootstrap_servers => "your-kafka-broker:9092" topics => ["Topic1"] codec => json } } filter { # Generate a unique fingerprint from the ID field fingerprint { source => ["id"] target => "[@metadata][fingerprint]" method => "SHA1" } # Check if the fingerprint exists in Redis redis { host => "your-redis-host" key => "logstash-dedup:%{[@metadata][fingerprint]}" field => "[@metadata][exists]" command => "GET" } # Drop duplicate messages if [@metadata][exists] { drop {} } else { # Store the fingerprint in Redis with a 1-day TTL redis { host => "your-redis-host" key => "logstash-dedup:%{[@metadata][fingerprint]}" value => "1" command => "SET" expire => 86400 } } } output { kafka { bootstrap_servers => "your-kafka-broker:9092" topic_id => "Topic2" codec => json } }
Key Notes
- For distributed Logstash setups, Redis is mandatory to share state across instances (otherwise each instance tracks its own duplicates independently).
- Adjust the
expirevalue to match your retention needs for deduplication (e.g., set to 604800 for 7 days).
内容的提问来源于stack exchange,提问作者Payam Heshmati

