No description
Find a file
Шевяк Кирилл Игоревич (АГЕНТ) a02d422cb4 Объединенный запрос на вытягивание 238410: Task 1274025: BE | Оптимизировать логи в либe Кафки
1274025 - Оптимизировать логи в либe Кафки

Связанные рабочие элементы: #1274019, #1274025
2026-07-24 17:47:48 +03:00
build/azure-pipelines Merged PR 221981: User Story 947697 | DevOps | Правки для версионирования 2026-05-07 15:54:07 +03:00
docs Объединенный запрос на вытягивание 230663: #1202961 Рефакторинг Outbox 2026-07-03 10:56:03 +03:00
requirements Merged PR 207252: User Story 1065765: Библиотека по работе с Kafka 1124425 Dev | Реализация библиотеки работы с Kafka 2026-03-19 09:05:08 +03:00
src Объединенный запрос на вытягивание 238410: Task 1274025: BE | Оптимизировать логи в либe Кафки 2026-07-24 17:47:48 +03:00
tests Merged PR 236549: #1260691 ТРП 5. [DEV Net] Разработка интеграции CПР. Обработка паттерна Inbox. 2026-07-22 15:45:34 +03:00
.editorconfig Merged PR 207252: User Story 1065765: Библиотека по работе с Kafka 1124425 Dev | Реализация библиотеки работы с Kafka 2026-03-19 09:05:08 +03:00
.gitignore Merged PR 236549: #1260691 ТРП 5. [DEV Net] Разработка интеграции CПР. Обработка паттерна Inbox. 2026-07-22 15:45:34 +03:00
CHANGELOG.md Объединенный запрос на вытягивание 238410: Task 1274025: BE | Оптимизировать логи в либe Кафки 2026-07-24 17:47:48 +03:00
Directory.Build.props Объединенный запрос на вытягивание 230663: #1202961 Рефакторинг Outbox 2026-07-03 10:56:03 +03:00
Directory.Packages.props Merged PR 236549: #1260691 ТРП 5. [DEV Net] Разработка интеграции CПР. Обработка паттерна Inbox. 2026-07-22 15:45:34 +03:00
Ingos.CoreServices.Kafka.sln Объединенный запрос на вытягивание 229114: Task 1203865: BE | Переезд на новую кафку. Поднятие версии. 2026-06-05 16:53:57 +03:00
nuget.config Merged PR 207252: User Story 1065765: Библиотека по работе с Kafka 1124425 Dev | Реализация библиотеки работы с Kafka 2026-03-19 09:05:08 +03:00
package-version-variables.yml Merged PR 236549: #1260691 ТРП 5. [DEV Net] Разработка интеграции CПР. Обработка паттерна Inbox. 2026-07-22 15:45:34 +03:00
README.md Merged PR 236549: #1260691 ТРП 5. [DEV Net] Разработка интеграции CПР. Обработка паттерна Inbox. 2026-07-22 15:45:34 +03:00

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)