【问题标题】:Is it possible to replicate kafka topics without alias prefix with MirrorMaker2是否可以使用 MirrorMaker2 复制没有别名前缀的 kafka 主题
【发布时间】:2019-12-18 10:58:26
【问题描述】:

我正在尝试在 2 个集群之间设置复制,但不希望更改主题名称。例如,如果我有一个名为“some_topic”的主题,它会自动复制到“cluster1.some_topic”,我很确定这可以完成,但还没有找到正确的配置来更改它

我当前的配置“mirrormaker2.properties”

# Sample MirrorMaker 2.0 top-level configuration file
# Run with ./bin/connect-mirror-maker.sh connect-mirror-maker.properties 

# specify any number of cluster aliases
clusters = cluster1, cluster2

# connection information for each cluster
cluster1.bootstrap.servers = host1:9092,host2:9092,host3:9092
cluster2.bootstrap.servers = rep_host1:9092,rep_host2:9092,rep_host3:9092

# enable and configure individual replication flows
cluster1->cluster2.enabled = true
cluster1->cluster2.topics = sometopic.*

# customize as needed
# replication.policy.separator = _
# sync.topic.acls.enabled = false
# emit.heartbeats.interval.seconds = 5

供参考:

【问题讨论】:

    标签: apache-kafka apache-kafka-mirrormaker


    【解决方案1】:

    要“禁用”主题前缀并同时正确镜像主题属性,我必须提供一个自定义的复制策略,该策略也覆盖topicSource 方法。否则非默认主题属性(例如,"cleanup.policy=compact")没有被镜像,即使在重新启动镜像生成器之后也是如此。

    这是对我有用的完整程序:

    1. 将以下自定义复制策略编译打包成 .jar 文件(完整源代码见here):
    public class PrefixlessReplicationPolicy extends DefaultReplicationPolicy {
    
      private static final Logger log = LoggerFactory.getLogger(PrefixlessReplicationPolicy.class);
    
      private String sourceClusterAlias;
    
      @Override
      public void configure(Map<String, ?> props) {
        super.configure(props);
        sourceClusterAlias = (String) props.get(MirrorConnectorConfig.SOURCE_CLUSTER_ALIAS);
        if (sourceClusterAlias == null) {
          String logMessage = String.format("Property %s not found", MirrorConnectorConfig.SOURCE_CLUSTER_ALIAS);
          log.error(logMessage);
          throw new RuntimeException(logMessage);
        }
      }
    
      @Override
      public String formatRemoteTopic(String sourceClusterAlias, String topic) {
        return topic;
      }
    
      @Override
      public String topicSource(String topic) {
        return topic == null ? null : sourceClusterAlias;
      }
    
      @Override
      public String upstreamTopic(String topic) {
        return null;
      }
    }
    
    1. 将.jar 复制到${KAFKA_HOME/libs 目录中
    2. 通过在${KAFKA_HOME}/config/mm2.properties 中设置replication.policy.class 属性,将Mirror Maker 2 配置为使用该复制策略:
      replication.policy.class=ch.mawileo.kafka.mm2.PrefixlessReplicationPolicy
    

    【讨论】:

    • 这对我来说是一个很好的答案,但我有一个问题:你为什么不直接更改方法 formatRemoteTopic 并在其他方法上调用 super ? (在我这边工作)
    • 嗨@Naremy,您所说的“其他方法”是指topicSourceupstreamTopic 吗?当我保持这些方法不变时,主题属性的镜像不起作用。例如,在我的 Kafka 设置中,默认情况下主题不会被压缩;当我创建一个我想要压缩的主题时(所以我为该主题设置了"cleanup.policy=compact" 属性),下游端的镜像主题仍然具有默认属性(即没有压缩)。主题属性的镜像是否适用于您的解决方案?
    • 在未来看起来这将是不必要的。截至 2021 年 7 月,已将 IdentityReplicationPolicy 添加到 MirrorMaker2。它似乎还不是发布的一部分。 github.com/apache/kafka/blob/trunk/connect/mirror-client/src/…
    【解决方案2】:

    我能够使用此设置删除前缀:

    "replication.policy.separator": ""
    "source.cluster.alias": "",
    "target.cluster.alias": "",
    

    如果您的情况需要别名设置,我知道您应该使用其他 replicationPolicy 类。默认使用 DefaultReplicationPolicy 类 (https://kafka.apache.org/24/javadoc/org/apache/kafka/connect/mirror/DefaultReplicationPolicy.html)

    【讨论】:

      【解决方案3】:

      我认为上面的答案是不恰当的。

      在 Mirror Maker 2.0 中,如果要保持主题不变,则必须实现 ReplicationPolicy。

      你可以参考DefaultReplicationPolicy.class,然后覆盖formatRemoteTopic(),之后你必须删除sourceClusterAlias + separator。最后在mm2.properties中配置replication.policy.class

      我定义了MigrationReplicationPolicy.class

      replication.policy.class = org.apache.kafka.connect.mirror.MigrationReplicationPolicy
      

      你应该看到MirrorClientConfig,class,我知道你会明白的

      【讨论】:

      • 您可以发布一个指向您创建的自定义MigrationReplicationPolicy.class 的链接吗?
      【解决方案4】:

      使用 Kafka ConfluentINC 连接器映像版本 5.4.2 管理推送复制 属性是:

      connector.class=org.apache.kafka.connect.mirror.MirrorSourceConnector
      target.cluster.alias= 
      replication.factor=3
      tasks.max=3
      topics=.*
      source.cluster.alias= 
      target.cluster.bootstrap.servers=<broker1>,<broker2>,<broker3>
      replication.policy.separator= 
      value.converter=org.apache.kafka.connect.converters.ByteArrayConverter
      source.cluster.bootstrap.servers=<broker1>,<broker2>,<broker3>
      key.converter=org.apache.kafka.connect.converters.ByteArrayConverter
      

      1) 在 3 个参数后面留空:source.cluster.alias、replication.policy.separator、target.cluster.alias。

      2) 在目标 Kafka 上设置此镜像连接器,而不是在源上(仅执行拉取)

      此外,您还可以使用 Conductor 或 Kafka Connector UI Landoop 图像 -landoop/kafka-connect-ui

      这仍处于测试场景中,但看起来很有希望。

      【讨论】:

        【解决方案5】:

        从Kafka 3.0.0开始,设置就够了

        replication.policy.class=org.apache.kafka.connect.mirror.IdentityReplicationPolicy
        

        此外,marcin-wieloch 的答案 https://stackoverflow.com/a/60619233/12008693 中的 PrefixlessReplicationPolicy 不再适用于 3.0.0 (NullPointerException)。

        【讨论】:

          【解决方案6】:

          我正在尝试在 2 个集群之间设置复制,但需要在两个集群中使用相同的主题名称,而无需在 connect-mirror-maker.properties 中提供别名。

          默认情况下,复制的主题会根据源集群别名重命名。

              Source --> Target
              topic-1 --> source.topic-1
          

          您可以通过在连接器属性文件下将以下属性设置为空白来避免重命名主题。默认情况下,replication.policy.separator 属性是一个句点,然后通过将其与 source.cluster.alias 一起设置为空白,目标主题将与源主题具有相同的名称。

          replication.policy.separator=
          source.cluster.alias=
          target.cluster.alias=
          

          【讨论】:

            猜你喜欢
            • 2020-05-30
            • 2021-03-27
            • 2013-08-16
            • 2016-10-12
            • 2023-03-10
            • 1970-01-01
            • 1970-01-01
            • 2017-02-22
            • 2018-02-02
            相关资源
            最近更新 更多