【问题标题】:Can we connect to/from a Kafka compacted topic with the Flink kafka Upsert connector?我们可以使用 Flink kafka Upsert 连接器连接到/从 Kafka 压缩主题吗?
【发布时间】:2021-04-15 04:43:07
【问题描述】:

感觉很明显,但我还是要问,因为我在文档中找不到明确的确认:

Flink 1.12 中可用的 Flink Table API upsert kafka connector 的语义与 Kafka 压缩主题的语义非常匹配:将流解释为变更日志,并使用 NULL 值作为墓碑来标记删除。

所以我的假设是可以使用它来消费和生产压缩主题,并且它可能正是为此而制作的,尽管它应该与非假设其内容确实是变更日志的压缩主题。但是我很惊讶在文档的那部分中没有找到任何对压缩主题的引用。

有人可以证实或证实这个假设吗?

【问题讨论】:

    标签: apache-kafka apache-flink flink-table-api


    【解决方案1】:

    是的,它是为压缩主题而设计的。根据FLIP-149

    一般来说,upsert-kafka源码的底层topic必须是compact的。此外,底层主题必须在同一个分区中具有相同key的所有数据,否则结果会出错。

    【讨论】:

    • 谢谢。 FLIP 中确实有更多信息,包括压缩主题和UPDATE_BEFORE 事件的行为。不幸的是,我看到的所有示例都只显示了连接器的接收端,尽管文档还描述了源的机制是我们所期望的,我想我现在明白了。
    • 您可以在this repository 中找到使用kafka-upsert 作为源连接器的示例。
    • 您可以使用常规的 kafka debezium avro 格式进行压缩吗? Upsert Kafka 似乎不支持文档所说的 debezium。
    • **在我的案例中特别作为来源。
    • @david-anderson 文档说:一般来说,upsert-kafka 源 must 的底层主题被压缩,但在我的测试中,non-compacted 对我来说也可以正常工作,不确定我错过了一些东西..
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-06-27
    • 2017-02-02
    • 2023-03-07
    • 2015-09-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多