Add schema registry support to Kafka module (#6190)

This commit introduces a schema registry functionality by adding interfaces and classes to manage schema registry definitions in the Elsa.Kafka module. It updates KafkaOptions to include schema registries and modifies classes to support schema registry configurations for producers and consumers. Additionally, it updates package references to include necessary dependencies for schema registry support.
This commit is contained in:
Sipke Schoorstra 2024-12-07 20:07:30 +01:00 committed by GitHub
parent a2448523a9
commit 4651b7a46d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
16 changed files with 113 additions and 20 deletions

View file

@ -15,7 +15,8 @@
<PackageVersion Include="BenchmarkDotNet" Version="0.14.0" />
<PackageVersion Include="Bogus" Version="35.6.1" />
<PackageVersion Include="ConfigureAwait.Fody" Version="3.3.2" PrivateAssets="All" />
<PackageVersion Include="Confluent.Kafka" Version="2.6.0" />
<PackageVersion Include="Confluent.Kafka" Version="2.6.1" />
<PackageVersion Include="Confluent.SchemaRegistry.Serdes.Avro" Version="2.4.0" />
<PackageVersion Include="coverlet.collector" Version="6.0.2" PrivateAssets="All" />
<PackageVersion Include="Cronos" Version="0.8.4" />
<PackageVersion Include="Dapper" Version="2.1.35" />

View file

@ -217,6 +217,15 @@
"Name": "topic-2"
}
],
"SchemaRegistries": [
{
"Id": "schema-registry-1",
"Name": "Schema Registry 1",
"Config": {
"Url": "http://localhost:8081"
}
}
],
"Producers": [
{
"Id": "producer-1",
@ -224,7 +233,8 @@
"FactoryType": "Elsa.Kafka.Factories.GenericProducerFactory`2[[System.String, System.Private.CoreLib], [Elsa.Server.Web.Messages.OrderReceived, Elsa.Server.Web]], Elsa.Kafka",
"Config": {
"BootstrapServers": "localhost:9092"
}
},
"SchemaRegistryId": "schema-registry-1"
}
],
"Consumers": [
@ -237,7 +247,8 @@
"GroupId": "group-1",
"AutoOffsetReset": "earliest",
"EnableAutoCommit": "true"
}
},
"SchemaRegistryId": "schema-registry-1"
}
]
},

View file

@ -5,10 +5,9 @@ using Elsa.Kafka.UIHints;
using Elsa.Workflows;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.UIHints;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Open.Linq.AsyncExtensions;
namespace Elsa.Kafka.Activities;
@ -73,7 +72,7 @@ public class ProduceMessage : CodeActivity
if (key is string keyString && string.IsNullOrWhiteSpace(keyString))
key = null;
using var producer = CreateProducer(context, producerDefinition);
using var producer = await CreateProducerAsync(context, producerDefinition);
var headers = CreateHeaders(context);
await producer.ProduceAsync(topic, key, content, headers, cancellationToken);
}
@ -94,14 +93,25 @@ public class ProduceMessage : CodeActivity
return headers;
}
private IProducer CreateProducer(ActivityExecutionContext context, ProducerDefinition producerDefinition)
private async Task<IProducer> CreateProducerAsync(ActivityExecutionContext context, ProducerDefinition producerDefinition)
{
var factory = context.GetOrCreateService(producerDefinition.FactoryType) as IProducerFactory;
if (factory == null)
throw new InvalidOperationException($"Producer factory of type '{producerDefinition.FactoryType}' not found.");
var createProducerContext = new CreateProducerContext(producerDefinition);
var schemaRegistryDefinition = await GetSchemaRegistryDefinitionAsync(context, producerDefinition.SchemaRegistryId);
var createProducerContext = new CreateProducerContext(producerDefinition, schemaRegistryDefinition);
return factory.CreateProducer(createProducerContext);
}
private async Task<SchemaRegistryDefinition?> GetSchemaRegistryDefinitionAsync(ActivityExecutionContext context, string? id, CancellationToken cancellationToken = default)
{
if (id == null)
return null;
var schemaRegistryDefinitionEnumerator = context.GetRequiredService<ISchemaRegistryDefinitionEnumerator>();
var schemaRegistryDefinitions = await schemaRegistryDefinitionEnumerator.EnumerateAsync(cancellationToken).ToList();
return schemaRegistryDefinitions.FirstOrDefault(x => x.Id == id);
}
}

View file

@ -0,0 +1,11 @@
namespace Elsa.Kafka;
public interface ISchemaRegistryDefinitionEnumerator
{
/// <summary>
/// Retrieves a list of schema registry definitions provided by <see cref="ISchemaRegistryDefinitionProvider"/> implementations.
/// </summary>
/// <param name="cancellationToken">A token to monitor for cancellation requests.</param>
/// <returns>A list of SchemaRegistryDefinition objects.</returns>
Task<IEnumerable<SchemaRegistryDefinition>> EnumerateAsync(CancellationToken cancellationToken);
}

View file

@ -0,0 +1,6 @@
namespace Elsa.Kafka;
public interface ISchemaRegistryDefinitionProvider
{
Task<IEnumerable<SchemaRegistryDefinition>> GetSchemaRegistryDefinitionsAsync(CancellationToken cancellationToken = default);
}

View file

@ -11,6 +11,7 @@
<PackageReference Include="Azure.Messaging.ServiceBus" />
<PackageReference Include="Azure.ResourceManager.ServiceBus" />
<PackageReference Include="Confluent.Kafka" />
<PackageReference Include="Confluent.SchemaRegistry.Serdes.Avro" />
<PackageReference Include="Microsoft.Extensions.Configuration.Abstractions" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" />
<PackageReference Include="System.Linq.Async" />

View file

@ -1,6 +1,5 @@
using Confluent.Kafka;
using Elsa.Common.Entities;
using Elsa.Kafka.Factories;
namespace Elsa.Kafka;
@ -9,4 +8,5 @@ public class ConsumerDefinition : Entity
public string Name { get; set; } = default!;
public Type FactoryType { get; set; } = default!;
public ConsumerConfig Config { get; set; } = new();
public string? SchemaRegistryId { get; set; } = default!;
}

View file

@ -1,6 +1,5 @@
using Confluent.Kafka;
using Elsa.Common.Entities;
using Elsa.Kafka.Factories;
namespace Elsa.Kafka;
@ -9,4 +8,5 @@ public class ProducerDefinition : Entity
public string Name { get; set; } = default!;
public Type FactoryType { get; set; } = default!;
public ProducerConfig Config { get; set; } = new();
public string? SchemaRegistryId { get; set; } = default!;
}

View file

@ -0,0 +1,10 @@
using Confluent.SchemaRegistry;
using Elsa.Common.Entities;
namespace Elsa.Kafka;
public class SchemaRegistryDefinition : Entity
{
public string Name { get; set; } = default!;
public SchemaRegistryConfig Config { get; set; } = default!;
}

View file

@ -60,12 +60,15 @@ public class KafkaFeature(IModule module) : FeatureBase(module)
.AddBackgroundTask<StartConsumersStartupTask>()
.AddSingleton<IWorkerManager, WorkerManager>()
.AddScoped<IWorkerTopicSubscriber, WorkerTopicSubscriber>()
.AddScoped<IConsumerDefinitionProvider, OptionsDefinitionProvider>()
.AddScoped<IProducerDefinitionProvider, OptionsDefinitionProvider>()
.AddScoped<ITopicDefinitionProvider, OptionsDefinitionProvider>()
.AddScoped<OptionsDefinitionProvider>()
.AddScoped<IConsumerDefinitionProvider>(sp => sp.GetRequiredService<OptionsDefinitionProvider>())
.AddScoped<IProducerDefinitionProvider>(sp => sp.GetRequiredService<OptionsDefinitionProvider>())
.AddScoped<ITopicDefinitionProvider>(sp => sp.GetRequiredService<OptionsDefinitionProvider>())
.AddScoped<ISchemaRegistryDefinitionProvider>(sp => sp.GetRequiredService<OptionsDefinitionProvider>())
.AddScoped<IConsumerDefinitionEnumerator, ConsumerDefinitionEnumerator>()
.AddScoped<IProducerDefinitionEnumerator, ProducerDefinitionEnumerator>()
.AddScoped<ITopicDefinitionEnumerator, TopicDefinitionEnumerator>()
.AddScoped<ISchemaRegistryDefinitionEnumerator, SchemaRegistryDefinitionEnumerator>()
.AddScoped<IPropertyUIHandler, ConsumerDefinitionsDropdownOptionsProvider>()
.AddScoped<IPropertyUIHandler, ProducerDefinitionsDropdownOptionsProvider>()
.AddScoped<IPropertyUIHandler, TopicDefinitionsDropdownOptionsProvider>()

View file

@ -0,0 +1,23 @@
using System.Runtime.CompilerServices;
using Open.Linq.AsyncExtensions;
namespace Elsa.Kafka.Implementations;
public class SchemaRegistryDefinitionEnumerator(IEnumerable<ISchemaRegistryDefinitionProvider> providers) : ISchemaRegistryDefinitionEnumerator
{
public async Task<IEnumerable<SchemaRegistryDefinition>> EnumerateAsync(CancellationToken cancellationToken)
{
return await GetSchemaRegistryDefinitionsInternalAsync(cancellationToken).ToListAsync(cancellationToken);
}
private async IAsyncEnumerable<SchemaRegistryDefinition> GetSchemaRegistryDefinitionsInternalAsync([EnumeratorCancellation] CancellationToken cancellationToken)
{
foreach (var provider in providers)
{
var schemaRegistryDefinitions = await provider.GetSchemaRegistryDefinitionsAsync(cancellationToken).ToList();
foreach (var schemaRegistryDefinition in schemaRegistryDefinitions)
yield return schemaRegistryDefinition;
}
}
}

View file

@ -41,7 +41,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) :
if (existingWorker == null)
{
// Add a new worker.
var worker = CreateWorker(scope.ServiceProvider, consumerDefinition);
var worker = await CreateWorkerAsync(scope.ServiceProvider, consumerDefinition, cancellationToken);
workers.Add(consumerDefinition.Id, worker);
}
else
@ -54,7 +54,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) :
if (existingHash != hash)
{
existingWorker.Stop();
var worker = CreateWorker(scope.ServiceProvider, consumerDefinition);
var worker = await CreateWorkerAsync(scope.ServiceProvider, consumerDefinition, cancellationToken);
workers[consumerDefinition.Id] = worker;
}
}
@ -162,7 +162,7 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) :
return Workers.TryGetValue(consumerDefinitionId, out var worker) ? worker : null;
}
private IWorker CreateWorker(IServiceProvider serviceProvider, ConsumerDefinition consumerDefinition)
private async Task<IWorker> CreateWorkerAsync(IServiceProvider serviceProvider, ConsumerDefinition consumerDefinition, CancellationToken cancellationToken)
{
var factoryType = consumerDefinition.FactoryType;
var consumerFactory = ActivatorUtilities.GetServiceOrCreateInstance(serviceProvider, factoryType) as IConsumerFactory;
@ -170,7 +170,8 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) :
if (consumerFactory == null)
throw new InvalidOperationException($"Worker factory of type '{factoryType}' not found.");
var createConsumerContext = new CreateConsumerContext(consumerDefinition);
var schemaRegistryDefinition = await GetSchemaRegistryDefinitionAsync(serviceProvider, consumerDefinition.SchemaRegistryId, cancellationToken);
var createConsumerContext = new CreateConsumerContext(consumerDefinition, schemaRegistryDefinition);
var consumerProxy = consumerFactory.CreateConsumer(createConsumerContext);
var wrappedConsumer = consumerProxy.Consumer;
var wrappedConsumerType = wrappedConsumer.GetType();
@ -186,4 +187,14 @@ public class WorkerManager(IHasher hasher, IServiceScopeFactory scopeFactory) :
{
return hasher.Hash(consumerDefinition);
}
private async Task<SchemaRegistryDefinition?> GetSchemaRegistryDefinitionAsync(IServiceProvider serviceProvider, string? id, CancellationToken cancellationToken = default)
{
if (id == null)
return null;
var schemaRegistryDefinitionEnumerator = serviceProvider.GetRequiredService<ISchemaRegistryDefinitionEnumerator>();
var schemaRegistryDefinitions = await schemaRegistryDefinitionEnumerator.EnumerateAsync(cancellationToken).ToList();
return schemaRegistryDefinitions.FirstOrDefault(x => x.Id == id);
}
}

View file

@ -1,3 +1,3 @@
namespace Elsa.Kafka;
public record CreateConsumerContext(ConsumerDefinition ConsumerDefinition);
public record CreateConsumerContext(ConsumerDefinition ConsumerDefinition, SchemaRegistryDefinition? SchemaRegistryDefinition);

View file

@ -1,3 +1,3 @@
namespace Elsa.Kafka;
public record CreateProducerContext(ProducerDefinition ProducerDefinition);
public record CreateProducerContext(ProducerDefinition ProducerDefinition, SchemaRegistryDefinition? SchemaRegistryDefinition);

View file

@ -3,6 +3,7 @@ namespace Elsa.Kafka;
public class KafkaOptions
{
public ICollection<TopicDefinition> Topics { get; set; } = [];
public ICollection<SchemaRegistryDefinition> SchemaRegistries { get; set; } = [];
public ICollection<ConsumerDefinition> Consumers { get; set; } = [];
public ICollection<ProducerDefinition> Producers { get; set; } = [];
public string WorkflowInstanceIdHeaderKey { get; set; } = "x-workflow-instance-id";

View file

@ -2,7 +2,7 @@ using Microsoft.Extensions.Options;
namespace Elsa.Kafka.Providers;
public class OptionsDefinitionProvider(IOptions<KafkaOptions> options) : IConsumerDefinitionProvider, IProducerDefinitionProvider, ITopicDefinitionProvider
public class OptionsDefinitionProvider(IOptions<KafkaOptions> options) : IConsumerDefinitionProvider, IProducerDefinitionProvider, ITopicDefinitionProvider, ISchemaRegistryDefinitionProvider
{
public Task<IEnumerable<ConsumerDefinition>> GetConsumerDefinitionsAsync(CancellationToken cancellationToken = default)
{
@ -18,4 +18,9 @@ public class OptionsDefinitionProvider(IOptions<KafkaOptions> options) : IConsum
{
return Task.FromResult<IEnumerable<TopicDefinition>>(options.Value.Topics);
}
public Task<IEnumerable<SchemaRegistryDefinition>> GetSchemaRegistryDefinitionsAsync(CancellationToken cancellationToken = default)
{
return Task.FromResult<IEnumerable<SchemaRegistryDefinition>>(options.Value.SchemaRegistries);
}
}