Apache Flink:如何根据事件类型将事件接收到不同的Kafka主题?

The*_*One 1 streaming apache-flink flink-streaming flink-cep

我想知道是否可以使用Flink Kafka接收器根据事件的类型编写不同主题的事件?假设我们有不同类型的事件:通知,消息和好友请求。我们希望将这些事件流式传输到不同的主题,这些主题分别是:notification-topic,messages-topic,friendsRequest-topic。

我尝试了多种方法来解决此问题,但仍然找不到正确的解决方案。我听说可以使用,ProcessFunction但如何将其与我的问题联系起来?

Ale*_*lex 5

如果您使用的是Kafka:

FlinkKafkaProducer011<Event> producer = new FlinkKafkaProducer011<>(
    "default.topic",

    new KeyedSerializationSchema<Event>() {
    @Override
    public byte[] serializeKey( Event element ) {
        return null; or element.getKey to bytes...
    }

    @Override
    public byte[] serializeValue( Event element ) {
            return event.toBytes() ...
    }

    @Override
    public String getTargetTopic( Event element ) {
        return element.getTopic();
    }
    },
    parameterTool.getProperties());

input.addSink(producer);
Run Code Online (Sandbox Code Playgroud)

它将getTargetTopic针对每个事件进行调用,以获取您要将事件路由到的主题。它将覆盖“ default.topic”