【问题标题】:Using kafka-streams with custom partitioner将 kafka-streams 与自定义分区器一起使用
【发布时间】:2021-07-14 02:14:52
【问题描述】:

我想通过 KTable 加入 KStream。两者都有不同的密钥,但使用自定义分区器共同分区。但是,连接不会产生结果。

KStream 具有以下结构
- 键:房屋 - 组
- 值:用户
KTable 具有以下结构
- 键:用户 - 组
- 值:地址

为了确保每个插入两个主题都按插入顺序进行处理,我使用了一个自定义分区器,我使用每个键的 Group 部分对两个主题进行分区。

我希望得到以下结构的流:
- 键:房屋 - 组
- 值:用户 - 地址

为此,我正在执行以下操作:

val streamsBuilder = streamBuilderHolder.streamsBuilder
val houseToUser = streamsBuilder.stream<HouseGroup, User>("houseToUser")
val userToAddress = streamsBuilder.table<UserGroup, Address>("userToAddress")
val result: KStream<HouseGroup, UserWithAddress> = houseToUser
        .map { k: HouseGroup, v: User ->
            val newKey = UserGroup(v, k.group)
            val newVal = UserHouse(v, k.house)
            KeyValue(newKey, newVal)
        }
        .join(userToAddress) { v1: UserHouse, v2: Address ->
            UserHouseWithAddress(v1, v2)
        }
        .map{k: UserGroup, v: UserHouseWithAddress ->
            val newKey = HouseGroup(v.house, k.group)
            val newVal = UserWithAddress(k.user, v.address)
            KeyValue(newKey, newVal)
        }

这需要一个匹配的连接,但它不起作用。

我想显而易见的解决方案是加入一个全局表并放弃自定义分区器。但是,我仍然不明白为什么上述方法不起作用。

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams


    【解决方案1】:

    我认为缺少匹配是因为使用了不同的分区器。

    对于您的输入主题,使用CustomPartitioner。 Kafka Streams 默认使用org.apache.kafka.clients.producer.internals.DefaultPartitioner

    KStream::join 之前的代码中,您调用了KStream::mapKStream::map 函数在 KStream::join 之前强制重新分区。在重新分区期间,消息被刷新到 Kafka($AppName-KSTREAM-MAP-000000000X-repartition 主题)。为了传播消息,Kafka Streams 使用定义的分区器(属性:ProducerConfig.PARTITIONER_CLASS_CONFIG)。总结:对于“repartition topic”和“KTable topic”,具有相同键的消息可能在不同的分区中

    您的解决方案将在您的 Kafka Streams 应用程序 (props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, "com.example.CustomPartitioner") 的属性中设置您的自定义分区

    对于调试,您可以检查重新分区主题 ($AppName-KSTREAM-MAP-000000000X-repartition)。具有相同键的消息(如输入主题)可能位于不同的分区(不同的编号)

    关于Join co-partitioning requirements的文档

    【讨论】:

    • 嗨 @wardziniak,您认为在 kafka 流中使用自定义分区器是一种好习惯吗?
    • @JanBols,这取决于您的分区器的 程度 - 消息的分布是否均匀跨分区。由于您已经使用自定义分区器(输入主题中的消息使用自定义分区器分发),我认为可以在您的用例中设置它。如果您想使用 Default one,则必须重新分区两个输入主题,因此效率不高。
    【解决方案2】:

    试试这个,对我有用。

    static async System.Threading.Tasks.Task Main(string[] args)
            {
     
                int count = 0;
                string line = null;
    
                var appConfig = getAppConfig(Enviroment.Dev);
                var schemaRegistrConfig = getSchemmaRegistryConfig(appConfig);
                var registry = new CachedSchemaRegistryClient(schemaRegistrConfig);
                var serializer = new AvroSerializer<YourAvroSchemaClass>(registry);
    
                var adminClient = new AdminClientBuilder(new AdminClientConfig( getClientConfig(appConfig))).Build();
                var topics = new List<TopicSpecification>(){ new TopicSpecification { Name = appConfig.OutputTopic, NumPartitions = 11}};
    
                await adminClient.CreateTopicsAsync(topics);
    
                var producerConfig = getProducerConfig(appConfig);
    
                var producer = new ProducerBuilder<string, byte[]>(producerConfig)
                    .SetPartitioner(appConfig.OutputTopic, (string topicName, int partitionCount, ReadOnlySpan<byte> keyData, bool keyIsNull) =>
                 {
                     var keyValueInInt = Convert.ToInt32(System.Text.UTF8Encoding.UTF8.GetString(keyData.ToArray()));
                     return (Partition)Math.Floor((double)(keyValueInInt % partitionCount));
                 }).Build();
    
                using (producer)
                {
                    Console.WriteLine($"Start to load data from : {appConfig.DataFileName}: { DateTime.Now} ");
                    var watch = new Stopwatch();
                    watch.Start();
                    try
                    {
                        var stream = new StreamReader(appConfig.DataFileName);
                        while ((line = stream.ReadLine()) != null)
                        {
                            var message = parseLine(line);
                            var data = await serializer.SerializeAsync(message.Value, new SerializationContext(MessageComponentType.Value, appConfig.OutputTopic));
                            producer.Produce(appConfig.OutputTopic, new Message<string, byte[]> { Key = message.Key, Value = data });
     
                            if (count++ % 1000 == 0)
                            {
                                producer.Flush();
                                Console.WriteLine($"Write ... {count} in {watch.Elapsed.TotalSeconds} seconds");
                            }
                        }
                        producer.Flush();
                    }
                    catch (ProduceException<Null, string> e)
                    {
                        Console.WriteLine($"Line: {line}");
                        Console.WriteLine($"Delivery failed: {e.Error.Reason}");
                        System.Environment.Exit(101);
                    }
            finally
            {
                producer.Flush();
               
                    }
                }
            }
    

    【讨论】:

      猜你喜欢
      • 2019-10-24
      • 2017-07-03
      • 2020-06-03
      • 1970-01-01
      • 1970-01-01
      • 2018-08-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多