【问题标题】:Send/produce json message through kafka通过 kafka 发送/生成 json 消息
【发布时间】:2022-01-02 09:12:13
【问题描述】:

这是我第一次使用 Kafka,我打算将 kafka 与 .net 一起使用

我想知道是否可以在生成事件时将 JSON 作为消息发送

我正在关注教程:https://developer.confluent.io/get-started/dotnet/#build-producer

另外,有没有办法将该值映射到模型,以便 value/json 结构始终与该模型绑定

例如:如果我希望我的 json 值为

{
  "customerName":"anything",
  "eventType":"one-of-three-enums",
  "columnsChanged": "string value or something"
}

我能找到的大部分例子都是这样的:

using Confluent.Kafka;
using System;
using Microsoft.Extensions.Configuration;

class Producer {
    static void Main(string[] args)
    {
        if (args.Length != 1) {
            Console.WriteLine("Please provide the configuration file path as a command line argument");
        }

        IConfiguration configuration = new ConfigurationBuilder()
            .AddIniFile(args[0])
            .Build();

        const string topic = "purchases";

        string[] users = { "eabara", "jsmith", "sgarcia", "jbernard", "htanaka", "awalther" };
        string[] items = { "book", "alarm clock", "t-shirts", "gift card", "batteries" };

        using (var producer = new ProducerBuilder<string, string>(
            configuration.AsEnumerable()).Build())
        {
            var numProduced = 0;
            const int numMessages = 10;
            for (int i = 0; i < numMessages; ++i)
            {
                Random rnd = new Random();
                var user = users[rnd.Next(users.Length)];
                var item = items[rnd.Next(items.Length)];

                producer.Produce(topic, new Message<string, string> { Key = user, Value = item },
                    (deliveryReport) =>
                    {
                        if (deliveryReport.Error.Code != ErrorCode.NoError) {
                            Console.WriteLine($"Failed to deliver message: {deliveryReport.Error.Reason}");
                        }
                        else {
                            Console.WriteLine($"Produced event to topic {topic}: key = {user,-10} value = {item}");
                            numProduced += 1;
                        }
                    });
            }

            producer.Flush(TimeSpan.FromSeconds(10));
            Console.WriteLine($"{numProduced} messages were produced to topic {topic}");
        }
    }
}

我希望该项目是 json 结构中的一个类。

【问题讨论】:

  • 您只需将var item 替换为作为字符串的序列化 JSON 对象...您应该能够在不了解 Kafka 的情况下执行此操作。另外,您是否在这里找到了示例代码? github.com/confluentinc/confluent-kafka-dotnet/blob/master/…
  • @OneCricketeer 如果事件不是相同的 json 结构,是否可以拒绝该事件。我希望它是强类型的?
  • 开箱即用,没有。这就是模式注册表的用途;存储消息必须匹配的 JSONSchema。或者,您可以使用 Protobuf 或其他强类型二进制格式
  • @OneCricketeer 只是一个后续问题。你知道我是否可以限制消息大小吗?有没有办法让生产者在发送之前检查事件的大小?

标签: .net asp.net-core apache-kafka


【解决方案1】:

想知道我是否可以在生成事件时将 JSON 作为消息发送

是的。 Kafka 存储字节并使用序列化器转换字节。在构建 Producer 时,您可以选择调用 SetValueSerializer

一些内置的序列化程序可以在 -https://github.com/confluentinc/confluent-kafka-dotnet/blob/master/src/Confluent.Kafka/Serializers.cs 找到

您需要自己编写才能通用地处理任何 JSON 模型类型。

将 Utf8Serializer 用于字符串时,您需要预先序列化模型类中的对象,然后将其作为值发送。在您的示例中,您将 var item 替换为一些序列化对象。

How do I turn a C# object into a JSON string in .NET?

使用模型类时,您的数据通常是强类型的,直到您开始手动编写 JSON 或使用 Dictionary 类型。如果您想要外部消息验证,Confluent Schema Registry 就是一个支持 JSONSchema 的示例,而来自 confluent-dotnet-kafka 项目的 JsonSerializer 支持这一点。

【讨论】:

  • 只是一个后续问题。你知道我是否可以限制消息大小吗?有没有办法让生产者在发送之前检查事件的大小,如果大小超过限制,则不允许发送消息?
  • Kafka 的默认消息批次限制为 1MB。如果你得到序列化字节数组的大小,那应该是单个记录大小的近似值(不过还有额外的开销,比如记录头和时间戳)
  • 谢谢。能否请您回答:stackoverflow.com/questions/70097676/…
猜你喜欢
  • 1970-01-01
  • 2018-11-01
  • 1970-01-01
  • 2014-04-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多