No description
- C# 100%
1274025 - Оптимизировать логи в либe Кафки Связанные рабочие элементы: #1274019, #1274025 |
||
|---|---|---|
| build/azure-pipelines | ||
| docs | ||
| requirements | ||
| src | ||
| tests | ||
| .editorconfig | ||
| .gitignore | ||
| CHANGELOG.md | ||
| Directory.Build.props | ||
| Directory.Packages.props | ||
| Ingos.CoreServices.Kafka.sln | ||
| nuget.config | ||
| package-version-variables.yml | ||
| README.md | ||
Ingos.CoreServices.Kafka
.NET 8 обёртка над Confluent.Kafka, снижающая boilerplate и добавляющая middleware pipeline, observability, failure strategies и гарантии доставки.
Пакеты
| Пакет | Назначение |
|---|---|
Ingos.CoreServices.Kafka |
Ядро: продюсер, консьюмер, middleware, конфигурация, стратегии ошибок, гарантии доставки |
Ingos.CoreServices.Kafka.Inbox.PostgreSql |
Exactly-once через idempotency key + inbox (Kafka → БД) — детали |
Ingos.CoreServices.Kafka.Outbox.PostgreSql |
Transactional Outbox паттерн (БД → Kafka) |
Ingos.CoreServices.Kafka.FeatureToggles |
Фиче-тогглы для продьюсеров и консьюмеров |
Quick Start
1. Установка
dotnet add package Ingos.CoreServices.Kafka
2. Регистрация в DI
services.AddOptions<Ingos.CoreServices.Kafka.Abstractions.KafkaOptions>()
.BindConfiguration(IngosKafkaOptions.SectionName);
services.AddOptions<KafkaConsumerOptions<OrderCreatedEvent>>()
.BindConfiguration("OrderCreatedEventConsumer");
var myConsumerOptions = configuration.GetRequiredSection("OrderCreatedEventConsumer")
.Get<KafkaConsumerOptions<OrderCreatedEvent>>()!;
services.AddKafka(kafka =>
{
// Продюсер
kafka.AddProducer<OrderCreatedEvent>("orders-topic",
producer =>
{
producer.WithDeliveryGuarantee(DeliveryGuarantee.AtLeastOnce);
producer.RetryPolicy(retry =>
{
retry.MaxRetries = 5;
retry.InitialDelay = TimeSpan.FromMilliseconds(200);
});
});
// Консьюмер
kafka.AddConsumer<OrderCreatedEvent, OrderCreatedHandler>(
myConsumerOptions.TopicName,
myConsumerOptions.GroupId,
consumer =>
{
consumer.WithDeliveryGuarantee(DeliveryGuarantee.AtLeastOnce);
consumer.RetryPolicy(retry =>
{
retry.MaxRetries = 3;
});
consumer.OnFailure(f => f
.ThenSendToDeadLetter<KafkaOptions>()
.AndAlert());
});
});
Пример настроек appsettings.json:
"IngosKafka": {
"SaslMechanism": "ScramSha512",
"SecurityProtocol": "SaslSsl",
"SaslUsername": "myKafkaUser",
"SaslPassword": "myKafkaPassword"
},
"OrderCreatedEventConsumer":
{
"PollTimeout": "00:00:00.3000000",
"TopicName": "main-queuing-osago-events-json",
"GroupId": "ingos-osago-events-group-dev"
}
3. Отправка сообщений
public class OrderService
{
private readonly IKafkaProducer<OrderCreatedEvent> _producer;
public OrderService(IKafkaProducer<OrderCreatedEvent> producer)
{
_producer = producer;
}
public async Task CreateOrder(Order order)
{
var @event = new OrderCreatedEvent { OrderId = order.Id };
// Простая отправка
await _producer.ProduceAsync(@event);
// С ключом и заголовками
await _producer.ProduceAsync(@event, ctx => ctx
.WithKey(order.Id.ToString())
.WithHeader("correlation-id", correlationId));
}
}
4. Обработка сообщений
public class OrderCreatedHandler : IMessageHandler<OrderCreatedEvent>
{
public async Task HandleAsync(
OrderCreatedEvent message,
ConsumeContext context,
CancellationToken ct)
{
// context.Topic, context.Partition, context.Offset
// context.Key, context.Headers, context.Timestamp
await ProcessOrderAsync(message, ct);
}
}
Опционально: подключение нескольких брокеров
При подключении нескольких брокеров необходимо для каждого зарегистрировать отдельный тип опций
var newKafkaOptions = services.SetupRequiredOptions<IngosKafkaOptionsNew>(configuration, $"{IngosKafkaOptionsNew.SectionName}");
private static TOptions SetupRequiredOptions<TOptions>(this IServiceCollection services,
IConfiguration configuration,
string optionsSection)
where TOptions : class
{
services.AddOptions<TOptions>().BindConfiguration(optionsSection);
return configuration.GetRequiredSection(optionsSection).Get<TOptions>()!;
}
а затем добавить каждый брокер отдельно.
При этом консьюмеры регистрируются как обычно AddConsumer,
а продьюсеров, при необходимости параллельной записи в несколько брокеров, регистрируют через AddClusterAwareKafkaProducer.
При необходимости регулировать запись/чтение в определенный брокер, используйте библиотеку Ingos.CoreServices.Kafka.FeatureToggles
services.AddKafka<IngosKafkaOptionsNew>(kafka =>
{
kafka.AddConsumer<MyType, AgreementsEventsIssuedHandler>(
myConsumerOptions.TopicName,
myConsumerOptions.GroupId,
consumer =>
{
consumer.UseFeatureToggle(myConsumerOptions); // Ingos.CoreServices.Kafka.FeatureToggles
});
kafka.AddClusterAwareKafkaProducer<MyType>( // AddClusterAwareKafkaProducer зарегистрирует продьюсера в пуле продьюсеров с одинаковым типом
myProducerOptions.TopicName,
producer =>
{
producer.UseFeatureToggle(myProducerOptions); // Ingos.CoreServices.Kafka.FeatureToggles
});
}
Возможности
- Типизированный продюсер —
IKafkaProducer<T>с перегрузками для ключей, заголовков, партиций (подробнее) - Типизированный консьюмер —
IMessageHandler<T>, пакетная обработка черезIBatchMessageHandler<T>(подробнее) - Гарантии доставки — AtMostOnce, AtLeastOnce из коробки (подробнее)
- Стратегии обработки ошибок — Dead Letter, Log & Skip, Alert, комбинирование (подробнее)
- Middleware pipeline — для продюсера и консьюмера (подробнее)
- Наблюдаемость — метрики через System.Diagnostics.Metrics, health checks (подробнее)
- Inbox / Outbox — Inbox: exactly-once на consumer; Outbox: at-least-once в Kafka +
x-idempotency-key(подробнее, примечания) | - Конфигурация — fluent API, appsettings.json, runtime reconfig (подробнее)
Структура решения
src/
├── Ingos.CoreServices.Kafka/
│ ├── Abstractions/ # Публичные интерфейсы и DTO
│ │ ├── Producer/ # IKafkaProducer, MessageContext, опции
│ │ ├── Consumer/ # IMessageHandler, ConsumeContext, опции
│ │ ├── Failure/ # IFailureStrategy, FailureContext
│ │ ├── Serialization/ # IMessageSerializer
│ │ ├── HeaderPropagation/ # Пропагация заголовков
│ │ └── Authorization/ # Контекст авторизации
│ ├── Configuration/ # DI-регистрация, fluent builders
│ └── Internal/ # Все реализации (internal)
├── Ingos.CoreServices.Kafka.Inbox.PostgreSql/
└── Ingos.CoreServices.Kafka.Outbox.PostgreSql/
Требования
- .NET 8+
- Confluent.Kafka 2.x
- PostgreSQL (для Inbox/Outbox)