If Apache Iggy is useful to you, give it a star on GitHubStar apache/iggy
Apache Iggy
SDKC#

Examples

Producer and consumer samples for the C# SDK, built on the high-level publisher and consumer.

These samples use the High-level SDK - the recommended way to build producers and consumers. For the low-level, per-call equivalents, see the Guide.

Use the installation and server setup first. The snippets share their named resources: run a producer before its corresponding consumer.

Producer

A publisher that creates the stream and topic if missing, batches sends in the background, and retries failures:

using System.Text;
using Apache.Iggy;
using Apache.Iggy.Configuration;
using Apache.Iggy.Enums;
using Apache.Iggy.Extensions;
using Apache.Iggy.Factory;
using Apache.Iggy.Messages;

using var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp
});

await client.ConnectAsync();
await client.LoginUserAsync("iggy", "iggy");

await using var publisher = client.CreatePublisherBuilder(
        Identifier.String("dev"),
        Identifier.String("events"))
    .CreateStreamIfNotExists("dev")
    .CreateTopicIfNotExists("events", topicPartitionsCount: 2)
    .WithBackgroundSending(batchSize: 100, flushInterval: TimeSpan.FromMilliseconds(100))
    .WithRetry(maxAttempts: 3)
    .Build();

await publisher.InitAsync();

for (var i = 0; i < 100; i++)
{
    var payload = Encoding.UTF8.GetBytes($"Event #{i}");
    await publisher.SendMessagesAsync(new List<Message> { new(Guid.NewGuid(), payload) });
}

// Drain the background queue before exiting
await publisher.WaitUntilAllSendsAsync();

Console.WriteLine("Queued 100 messages");

Consumer group

A consumer that creates and joins a consumer group, stores each offset after handling the yielded message and requesting the next one, and surfaces polling errors:

using System.Text;
using Apache.Iggy;
using Apache.Iggy.Configuration;
using Apache.Iggy.Consumers;
using Apache.Iggy.Enums;
using Apache.Iggy.Extensions;
using Apache.Iggy.Factory;
using Apache.Iggy.Kinds;

using var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp
});

await client.ConnectAsync();
await client.LoginUserAsync("iggy", "iggy");

await using var consumer = client.CreateConsumerBuilder(
        Identifier.String("dev"),
        Identifier.String("events"),
        Consumer.Group("event-processors"))
    .WithConsumerGroup("event-processors", createIfNotExists: true, joinGroup: true)
    .WithPollingStrategy(PollingStrategy.Next())
    .WithBatchSize(20)
    .WithAutoCommitMode(AutoCommitMode.AfterReceive)
    .SubscribeOnPollingError(e =>
    {
        Console.WriteLine($"Polling error: {e.Exception.Message}");
        return Task.CompletedTask;
    })
    .Build();

await consumer.InitAsync();

await foreach (var message in consumer.ReceiveAsync())
{
    var payload = Encoding.UTF8.GetString(message.Message.Payload);
    Console.WriteLine($"Partition {message.PartitionId}, offset {message.CurrentOffset}: {payload}");
}

Run several instances of the consumer to see the group load-balance partitions between members.

Typed messages

A typed publisher/consumer pair that (de)serializes a record as JSON. The typed builders are configured via statements rather than one fluent chain - the With* methods return the untyped base builder, so Build() must be called on the typed builder variable (see Typed consumer):

using System.Buffers;
using System.Text.Json;
using Apache.Iggy;
using Apache.Iggy.Configuration;
using Apache.Iggy.Consumers;
using Apache.Iggy.Enums;
using Apache.Iggy.Factory;
using Apache.Iggy.Kinds;
using Apache.Iggy.Publishers;

using var client = IggyClientFactory.CreateClient(new IggyClientConfigurator
{
    BaseAddress = "127.0.0.1:8090",
    Protocol = Protocol.Tcp
});

await client.ConnectAsync();
await client.LoginUserAsync("iggy", "iggy");

// Publish typed messages
var publisherBuilder = IggyPublisherBuilder<OrderEvent>.Create(
    client,
    Identifier.String("orders"),
    Identifier.String("created"),
    new OrderSerializer()
);
publisherBuilder.CreateStreamIfNotExists("orders");
publisherBuilder.CreateTopicIfNotExists("created");

await using var publisher = publisherBuilder.Build();
await publisher.InitAsync();

await publisher.SendAsync(new OrderEvent(Guid.NewGuid(), 99.90m));

// Consume them
var consumerBuilder = IggyConsumerBuilder<OrderEvent>.Create(
    client,
    Identifier.String("orders"),
    Identifier.String("created"),
    Consumer.New(1),
    new OrderDeserializer()
);
consumerBuilder.WithPollingStrategy(PollingStrategy.Next());
consumerBuilder.WithAutoCommitMode(AutoCommitMode.AfterReceive);

await using var consumer = consumerBuilder.Build();
await consumer.InitAsync();

await foreach (var message in consumer.ReceiveDeserializedAsync())
{
    if (message.Status == MessageStatus.Success)
    {
        Console.WriteLine($"Order {message.Data!.OrderId}, amount {message.Data.Amount}");
    }
}

record OrderEvent(Guid OrderId, decimal Amount);

class OrderSerializer : ISerializer<OrderEvent>
{
    public void Serialize(OrderEvent data, IBufferWriter<byte> writer) =>
        writer.Write(JsonSerializer.SerializeToUtf8Bytes(data));
}

class OrderDeserializer : IDeserializer<OrderEvent>
{
    public OrderEvent Deserialize(ReadOnlyMemory<byte> data) =>
        JsonSerializer.Deserialize<OrderEvent>(data.Span)!;
}

More examples

The examples/csharp directory in the Iggy repository contains complete, runnable projects:

  • Basic / GettingStarted - low-level producer and consumer
  • NewSdk - high-level IggyPublisher / IggyConsumer (like the samples above)
  • MessageEnvelope - envelope pattern (message type + JSON payload) over the low-level client
  • MessageHeaders - user-defined message headers
  • TcpTls - TLS-encrypted TCP connection

Repository examples require .NET 10 and reference the local SDK. Build and run from the same checkout as the server:

dotnet build examples/csharp/Iggy_SDK.Examples.sln --configuration Release
dotnet run --no-build --configuration Release --project examples/csharp/src/GettingStarted/Iggy_SDK.Examples.GettingStarted.Producer
dotnet run --no-build --configuration Release --project examples/csharp/src/GettingStarted/Iggy_SDK.Examples.GettingStarted.Consumer

The getting-started producer sends 50 messages to partition 0; its consumer reads five nonempty batches and then exits. Envelope and header examples also use five batches. NewSdk uses four partitions and its group consumer runs until stopped.

For TcpTls, run from the repository root so core/certs/iggy_ca_cert.pem resolves. Start a separate matching server with the example TLS certificate and root credentials:

IGGY_ROOT_USERNAME=iggy IGGY_ROOT_PASSWORD=iggy \
IGGY_TCP_TLS_ENABLED=true \
IGGY_TCP_TLS_CERT_FILE=core/certs/iggy_cert.pem \
IGGY_TCP_TLS_KEY_FILE=core/certs/iggy_key.pem \
cargo run --bin iggy-server
dotnet run --no-build --configuration Release --project examples/csharp/src/TcpTls/Iggy_SDK.Examples.TcpTls.Producer
dotnet run --no-build --configuration Release --project examples/csharp/src/TcpTls/Iggy_SDK.Examples.TcpTls.Consumer

On this page