Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions dotnet-client-libraries/RabbitMqVsKafka/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
# RabbitMQ vs Kafka for .NET Developers

The sample for the Code Maze article [RabbitMQ vs Kafka for .NET Developers](https://code-maze.com/rabbitmq-vs-kafka-dotnet/).

## Run the demos

Start both brokers from this folder (Docker required):

```bash
docker compose up -d
```

Then run one demo at a time from the same folder:

```bash
dotnet run --project RabbitMqVsKafka -- rabbit
dotnet run --project RabbitMqVsKafka -- kafka
dotnet run --project RabbitMqVsKafka -- stream
```

The demos leave messages, topics and committed offsets behind, so reset both brokers before you run a demo a second time:

```bash
docker compose down -v
docker compose up -d
```

## Run the tests

The tests start their own brokers with Testcontainers, so they need Docker but not the Compose file:

```bash
dotnet test
```
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
using Testcontainers.Kafka;
using Testcontainers.RabbitMq;

[assembly: CaptureConsole]

namespace RabbitMqVsKafka.Tests;

public class BrokerTests
{
private static readonly int[] GoodOrders = [1001, 1002, 1004, 1005];
private static readonly int[] AllOrders = [1001, 1002, 1003, 1004, 1005];

[Fact]
public async Task RabbitQueue_AcksGoodOrders_DeadLettersPoisonOrder_LeavesNothingForLateConsumer()
{
await using var rabbitMq = new RabbitMqBuilder("rabbitmq:4.3.6-management").Build();
await rabbitMq.StartAsync(TestContext.Current.CancellationToken);

var result = await new RabbitQueueDemo(rabbitMq.GetConnectionString()).RunAsync();

Assert.Equal(GoodOrders, result.Billed.Order());
Assert.Equal(1003, result.DeadLetteredOrderId);
Assert.Equal("delivery_limit", result.DeadLetterReason);
Assert.Equal(3, result.PoisonAttempts);
Assert.False(result.LateConsumerGotAnything);
}

[Fact]
public async Task Kafka_CommitsGoodOrders_DeadLettersPoisonOrder_SecondGroupAndReplayReadAllFive()
{
await using var kafka = new KafkaBuilder("apache/kafka:4.3.1")
.WithCommand(StartKafkaWithoutTrailingComma)
.Build();
await kafka.StartAsync(TestContext.Current.CancellationToken);

var result = await new KafkaDemo(kafka.GetBootstrapAddress()).RunAsync();

Assert.Equal(GoodOrders, result.Billed.Order());
Assert.Equal([1003], result.DeadLettered);
Assert.Equal(AllOrders, result.Analytics.Order());
Assert.Equal(AllOrders, result.Replayed.Order());
}

[Fact]
public async Task RabbitStream_TwoReadersFromFirstOffset_BothReadAllFive()
{
await using var rabbitMq = new RabbitMqBuilder("rabbitmq:4.3.6-management").Build();
await rabbitMq.StartAsync(TestContext.Current.CancellationToken);

var (first, second) = await new RabbitStreamDemo(rabbitMq.GetConnectionString()).RunAsync();

Assert.Equal(AllOrders, first);
Assert.Equal(AllOrders, second);
}

// Testcontainers.Kafka 4.15.0 writes KAFKA_ADVERTISED_LISTENERS with a trailing comma, and
// apache/kafka:4.3.1 refuses to start with it ("values must not be empty"). The fix
// (testcontainers-dotnet PR 1772) is not released yet, so this start command drops the
// comma from the generated startup script before it starts the broker.
private static readonly DotNet.Testcontainers.Configurations.OverwriteEnumerable<string> StartKafkaWithoutTrailingComma = new(
[
"while [ ! -f /testcontainers.sh ]; do sleep 0.1; done; " +
"sed 's/,$//' /testcontainers.sh > /tmp/testcontainers.sh && exec bash /tmp/testcontainers.sh"
]);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFramework>net10.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
<OutputType>Exe</OutputType>
<IsPackable>false</IsPackable>
<!-- Print each test's console output in the dotnet test log, so the CI log shows what the demos print. -->
<VSTestLogger>console%3Bverbosity=detailed</VSTestLogger>
</PropertyGroup>

<ItemGroup>
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="18.10.1" />
<PackageReference Include="Testcontainers.Kafka" Version="4.15.0" />
<PackageReference Include="Testcontainers.RabbitMq" Version="4.15.0" />
<PackageReference Include="xunit.v3.mtp-off" Version="4.0.1" />
<PackageReference Include="xunit.runner.visualstudio" Version="4.0.0" />
</ItemGroup>

<ItemGroup>
<Using Include="Xunit" />
</ItemGroup>

<ItemGroup>
<ProjectReference Include="..\RabbitMqVsKafka\RabbitMqVsKafka.csproj" />
</ItemGroup>

</Project>
48 changes: 48 additions & 0 deletions dotnet-client-libraries/RabbitMqVsKafka/RabbitMqVsKafka.sln
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@

Microsoft Visual Studio Solution File, Format Version 12.00
# Visual Studio Version 17
VisualStudioVersion = 17.0.31903.59
MinimumVisualStudioVersion = 10.0.40219.1
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "RabbitMqVsKafka", "RabbitMqVsKafka\RabbitMqVsKafka.csproj", "{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "RabbitMqVsKafka.Tests", "RabbitMqVsKafka.Tests\RabbitMqVsKafka.Tests.csproj", "{3C70F97A-C5EA-4F0C-A14C-1166377969AF}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Debug|x64 = Debug|x64
Debug|x86 = Debug|x86
Release|Any CPU = Release|Any CPU
Release|x64 = Release|x64
Release|x86 = Release|x86
EndGlobalSection
GlobalSection(ProjectConfigurationPlatforms) = postSolution
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|Any CPU.Build.0 = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|x64.ActiveCfg = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|x64.Build.0 = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|x86.ActiveCfg = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Debug|x86.Build.0 = Debug|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|Any CPU.ActiveCfg = Release|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|Any CPU.Build.0 = Release|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|x64.ActiveCfg = Release|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|x64.Build.0 = Release|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|x86.ActiveCfg = Release|Any CPU
{843910BC-21F7-4CBD-9DAB-D2E64EC8F380}.Release|x86.Build.0 = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|Any CPU.Build.0 = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|x64.ActiveCfg = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|x64.Build.0 = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|x86.ActiveCfg = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Debug|x86.Build.0 = Debug|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|Any CPU.ActiveCfg = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|Any CPU.Build.0 = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|x64.ActiveCfg = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|x64.Build.0 = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|x86.ActiveCfg = Release|Any CPU
{3C70F97A-C5EA-4F0C-A14C-1166377969AF}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
EndGlobalSection
EndGlobal
126 changes: 126 additions & 0 deletions dotnet-client-libraries/RabbitMqVsKafka/RabbitMqVsKafka/KafkaDemo.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
using System.Text.Json;
using Confluent.Kafka;
using Confluent.Kafka.Admin;

namespace RabbitMqVsKafka;

public record KafkaResult(
IReadOnlyList<int> Billed, IReadOnlyList<int> DeadLettered, IReadOnlyList<int> Analytics, IReadOnlyList<int> Replayed);

public class KafkaDemo(string bootstrapServers)
{
private const int MaxAttempts = 3;

public async Task<KafkaResult> RunAsync()
{
await CreateTopicsAsync();
var startedAt = DateTime.UtcNow;

using var producer = new ProducerBuilder<string, string>(
new ProducerConfig { BootstrapServers = bootstrapServers }).Build();
foreach (var order in OrderPlaced.Samples)
{
var sent = await producer.ProduceAsync("orders", new Message<string, string>
{
Key = order.OrderId.ToString(),
Value = JsonSerializer.Serialize(order)
});
Console.WriteLine($"Kafka: produced order {order.OrderId} to partition {sent.Partition.Value}, offset {sent.Offset.Value}");
}

var (billed, deadLettered) = await BillAsync(producer);
var analytics = Read("analytics", consumer => consumer.Subscribe("orders"));
var replayed = Read("billing-replay", consumer =>
{
var fromTime = Enumerable.Range(0, 3)
.Select(partition => new TopicPartitionTimestamp("orders", partition, new Timestamp(startedAt)));
consumer.Assign(consumer.OffsetsForTimes(fromTime, TimeSpan.FromSeconds(10)));
});

return new KafkaResult(billed, deadLettered, analytics, replayed);
}

private async Task<(List<int> Billed, List<int> DeadLettered)> BillAsync(IProducer<string, string> producer)
{
var billed = new List<int>();
var deadLettered = new List<int>();
var attempts = new Dictionary<int, int>();

using var consumer = CreateConsumer("billing");
consumer.Subscribe("orders");

while (billed.Count + deadLettered.Count < OrderPlaced.Samples.Count)
{
var record = consumer.Consume(TimeSpan.FromSeconds(30))!;
var order = JsonSerializer.Deserialize<OrderPlaced>(record.Message.Value)!;
var where = $"partition {record.Partition.Value}, offset {record.Offset.Value}";

if (order.CanBeBilled)
{
consumer.Commit(record);
billed.Add(order.OrderId);
Console.WriteLine($"Kafka billing: billed order {order.OrderId} ({where}), commit");
continue;
}

attempts[order.OrderId] = attempts.GetValueOrDefault(order.OrderId) + 1;
if (attempts[order.OrderId] < MaxAttempts)
{
Console.WriteLine($"Kafka billing: order {order.OrderId} failed ({where}), attempt {attempts[order.OrderId]}, seek back");
consumer.Seek(record.TopicPartitionOffset);
continue;
}

await producer.ProduceAsync("orders.dlq", record.Message);
consumer.Commit(record);
deadLettered.Add(order.OrderId);
Console.WriteLine($"Kafka billing: order {order.OrderId} failed {MaxAttempts} times, produced to orders.dlq, commit");
}

consumer.Close();
return (billed, deadLettered);
}

private List<int> Read(string groupId, Action<IConsumer<string, string>> start)
{
var received = new List<int>();
using var consumer = CreateConsumer(groupId);
start(consumer);

while (received.Count < OrderPlaced.Samples.Count)
{
var record = consumer.Consume(TimeSpan.FromSeconds(30))!;
var order = JsonSerializer.Deserialize<OrderPlaced>(record.Message.Value)!;
received.Add(order.OrderId);
Console.WriteLine($"Kafka {groupId}: read order {order.OrderId} (partition {record.Partition.Value}, offset {record.Offset.Value})");
}

consumer.Close();
return received;
}

private IConsumer<string, string> CreateConsumer(string groupId) =>
new ConsumerBuilder<string, string>(new ConsumerConfig
{
BootstrapServers = bootstrapServers,
GroupId = groupId,
AutoOffsetReset = AutoOffsetReset.Earliest,
EnableAutoCommit = false
}).Build();

private async Task CreateTopicsAsync()
{
using var admin = new AdminClientBuilder(new AdminClientConfig { BootstrapServers = bootstrapServers }).Build();
try
{
await admin.CreateTopicsAsync(
[
new TopicSpecification { Name = "orders", NumPartitions = 3, ReplicationFactor = 1 },
new TopicSpecification { Name = "orders.dlq", NumPartitions = 1, ReplicationFactor = 1 }
]);
}
catch (CreateTopicsException e) when (e.Results.All(r => r.Error.Code == ErrorCode.TopicAlreadyExists))
{
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
namespace RabbitMqVsKafka;

public record OrderPlaced(int OrderId, string Customer, decimal Total)
{
public static IReadOnlyList<OrderPlaced> Samples { get; } =
[
new(1001, "Ana", 49.90m),
new(1002, "Ben", 120.00m),
new(1003, "Chloe", -15.00m),
new(1004, "Dan", 75.50m),
new(1005, "Eva", 9.99m)
];

public bool CanBeBilled => Total > 0;
}
20 changes: 20 additions & 0 deletions dotnet-client-libraries/RabbitMqVsKafka/RabbitMqVsKafka/Program.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
using RabbitMqVsKafka;

const string rabbitMq = "amqp://guest:guest@localhost:5672";
const string kafka = "localhost:9092";

switch (args.FirstOrDefault())
{
case "rabbit":
await new RabbitQueueDemo(rabbitMq).RunAsync();
break;
case "kafka":
await new KafkaDemo(kafka).RunAsync();
break;
case "stream":
await new RabbitStreamDemo(rabbitMq).RunAsync();
break;
default:
Console.WriteLine("Usage: dotnet run -- rabbit | kafka | stream");
break;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net10.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>

<ItemGroup>
<PackageReference Include="Confluent.Kafka" Version="2.15.1" />
<PackageReference Include="RabbitMQ.Client" Version="7.2.2" />
</ItemGroup>

</Project>
Loading
Loading