【问题标题】:Kafka - Redirect messages from "Topic A" to "Topic B" based on header valueKafka - 根据标头值将消息从“主题 A”重定向到“主题 B”
【发布时间】:2019-03-27 16:43:50
【问题描述】:

我想将 kafka 消息从名为“all-topic”的主题重定向到名为“headervalue-topic”的主题,其中 headervalue 是自定义的值每条消息都有标题。

目前我正在使用一个自定义控制台应用程序,它使用消息并将消息重定向到正确的主题,但它每秒只处理 16 条消息。

kafka 和 zookeeper 都在 docker 容器中运行,配置如下:

zookeeper:
  image: "wurstmeister/zookeeper:latest"
  restart: always
  ports:
    - "2181:2181"
  environment:
    ZOOKEEPER_CLIENT_PORT: 2181
    ZOOKEEPER_SERVER_ID: 1

kafka:
  hostname: kafka
  image: "wurstmeister/kafka:latest"
  restart: always
  depends_on:
    - zookeeper
  ports:
    - "9092:9092"
  environment:
    KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
    KAFKA_ADVERTISED_HOST_NAME: kafka
    KAFKA_ADVERTISED_PORT: 9092

实现我的目标最好、最快的方法是什么?

我确实知道 Kafka Streams 的存在,但我对 Java 不熟悉,所以如果您想建议 Kafka Streams 一个小例子,我们将不胜感激:)

非常感谢!

【问题讨论】:

标签: apache-kafka apache-kafka-streams


【解决方案1】:

这是我想出的解决方案,使用 kafka-streams nodejs 库:

const {KafkaStreams} = require("kafka-streams");
const {nativeConfig: config} = require("./config.js");

const kafkaStreams = new KafkaStreams(config);
const myConsumerStream = kafkaStreams.getKStream("all-topic");

myConsumerStream
    .mapJSONConvenience()
    .filter((element) => {
        return element.value.type == "Article";
    })
    .tap((element) => {console.log("Got Article")})
    .mapWrapKafkaValue()
    .to("Article-topic", 1, "buffer");

myConsumerStream.start();

【讨论】:

    【解决方案2】:

    据我所知,您无法通过 DSL 直接访问标头。 不过,您可以使用流处理器通过 ProcessorContext 访问它,这是我想出的一个小例子:

    public class CustomProcessor1 implements Processor<String, String> {
    private ProcessorContext context;
    
    @Override
    public void init(ProcessorContext processorContext) {
        this.context = processorContext;
    }
    
    @Override
    public void process(String key, String value) {
        HashMap<String, String> headers = new HashMap<>();
        for (Header header : context.headers()) {
                headers.put(header.key(), new String(header.value()));
        }
        String headerValue = headers.get("certainHeader").replace("\"", "");
    
        if (headerValue.equals("expectedHeaderValue")) {
            context.forward(key, value);
        }
    }
    

    上面是处理器,它将带有与 headerValue 匹配的某些Header 的消息转发到下游进程。创建流拓扑时将使用此处理器,如下所示:

    public static void main(String[] args) throws Exception {
        Properties props = getProperties();
    
        final Topology topology = new Topology()
                .addSource("SOURCE", "all.topic")
                .addProcessor("CUSTOM_PROCESSOR_1", CustomProcessor1::new, "SOURCE")
                .addProcessor("CUSTOM_PROCESSOR_2", CustomProcessor2::new, "SOURCE")
                .addSink("SINK1", "headervalue1-topic", "CUSTOM_PROCESSOR_1")
                .addSink("SINK2", "headervalue2-topic", "CUSTOM_PROCESSOR_2");
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-03-10
      相关资源
      最近更新 更多