【发布时间】: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