【问题标题】:Apache Flink: How to sink events to different Kafka topics depending on the event type?Apache Flink:如何根据事件类型将事件下沉到不同的 Kafka 主题?
【发布时间】:2018-03-27 08:59:02
【问题描述】:
我想知道是否可以使用 Flink Kafka sink 根据事件类型编写不同主题的事件?
假设我们有不同类型的事件:通知、消息和好友请求。我们希望将这些事件流式传输到不同的主题:notification-topic、messages-topic、friendsRequest-topic。
我尝试了许多不同的方法来解决这个问题,但仍然找不到正确的解决方案。我听说我可以使用ProcessFunction,但这与我的问题有什么关系?
【问题讨论】:
标签:
streaming
apache-flink
flink-streaming
flink-cep
【解决方案1】:
如果您使用的是 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);
它将为每个事件调用getTargetTopic,以获取您想要将事件路由到的主题。它将覆盖“default.topic”