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

如何利用Logstash实现Kafka Topic消息去重并转发至新Topic?

Kafka Message Deduplication: Forward Unique Messages from Topic1 to Topic2

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.

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 expire value to match your retention needs for deduplication (e.g., set to 604800 for 7 days).

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:53:10