【发布时间】:2020-08-21 19:42:47
【问题描述】:
我最近使用 Spring Cloud Stream 和 Kafka 创建了一个微服务。我是两个世界的新手,所以如果我问了一个愚蠢的问题,请原谅我。
此服务的流程是从一个主题消费,转换数据,然后将结果生成到各个主题。
进入消费者主题的数据是数据库更改。基本上,我们有另一个服务监控数据库日志并产生对我的新服务正在消费的主题的更改。
在我的新服务中,我定义了消费者和生产者绑定。数据库数据以byte[]格式传入,消费者读取数据并将byte[]解码为Java对象DBData,然后根据表名进行不同的转换。
请看下面的代码示例。
@StreamListener
@SendTo("OUTPUT_MOCK")
public KStream<String, Mock> process(@Input("DB_SOURCE") KStream<String, byte[]> input) {
return input
.map(((String key, byte[] bytes) -> {
try {
return new KeyValue<>(key, DBDecoder.decode(bytes)); // decodes byte[] into DBData
} catch (Exception e) {
return new KeyValue<>(key, null);
}
}))
.filter((key, v) -> (v instanceof DBData)) // filter value type
.map((key, v) -> new KeyValue<>(key, (DBData) v))
.filter((String key, DBData v) -> v.getTableName().equalse("MOCK")) // check the table name
.flatMap((String key, DBData v) -> extractMockDataChanges(v)); // convert to Mock object
}
从代码示例中,您可以看到 DB 数据进入,然后被解码为 DBData 格式。然后根据表名过滤结果,最终生成转换后的Mock对象到OUTPUT_MOCK主题。
这段代码运行良好,但我的主要问题是DBData 转换部分。这个OUTPUT_MOCK 主题是许多其他生产者主题之一。我必须对许多其他表执行此操作,并且每次我都必须重复解码过程,这似乎是不必要且多余的。
是否有更好的方法来处理数据库数据转换,以便转换后的DBData 可用于其他流处理器?
PS:我查看了状态存储,但这似乎有点矫枉过正,因为 Kafka 在将数据添加到存储中时会对其进行序列化,而在提取时它将对它们进行反序列化。所以我试图避免额外的开销。
【问题讨论】:
-
您是否尝试使用
KStreamAPI 中的branch方法并使用谓词检查表名。这样,您可以为每个表设置一个谓词,然后相应地重定向到输出主题。
标签: apache-kafka spring-cloud-stream