High-level SDK
The per-call IIggyClient API is explicit but verbose — you manage partitioning, batching, retries, and offsets yourself. The high-level IggyPublisher and IggyConsumer wrap that surface with fluent builders that handle:
- Background sending with buffering, batching, and retries
- Automatic stream/topic/consumer-group creation
- Automatic offset commits (on each poll, after each received message, or manual)
- Consuming a topic as an async stream (
IAsyncEnumerable) - Typed message (de)serialization
- Sending and receiving over rented, pooled buffers — no per-message allocation
Both are built from an existing connected client via IggyPublisherBuilder.Create(...) / IggyConsumerBuilder.Create(...), or equivalently via the client.CreatePublisherBuilder(...) / client.CreateConsumerBuilder(...) extension methods from Apache.Iggy.Extensions.
IggyPublisher
Configure a publisher, initialize it, then send. InitAsync validates (and optionally creates) the stream and topic:
using Apache.Iggy;
using Apache.Iggy.Kinds;
using Apache.Iggy.Messages;
using Apache.Iggy.Publishers;
var publisher = IggyPublisherBuilder.Create(
client,
Identifier.String("my-stream"),
Identifier.String("my-topic")
)
.WithPartitioning(Partitioning.None())
.WithBackgroundSending(enabled: true, batchSize: 100)
.WithRetry(maxAttempts: 3)
.Build();
await publisher.InitAsync();
var messages = new List<Message>
{
new(Guid.NewGuid(), "Message 1"u8.ToArray()),
new(0, "Message 2"u8.ToArray())
};
await publisher.SendMessagesAsync(messages);
// Drain the background queue, then dispose
await publisher.WaitUntilAllSendsAsync();
await publisher.DisposeAsync();Publisher builder options
| Method | Description |
|---|---|
WithConnection(protocol, address, login, password, receiveBufferSize = 4096, sendBufferSize = 4096, reconnectionSettings = null) | Connection settings — only used when the builder creates its own client (i.e. the Create(streamId, topicId) overload, without an existing client) |
WithPartitioning(partitioning) | Routing strategy for produced messages (default: balanced) |
CreateStreamIfNotExists(name) | Auto-create the stream on InitAsync if missing |
CreateTopicIfNotExists(name, topicPartitionsCount = 1, compressionAlgorithm = None, replicationFactor = null, messageExpiry = TimeSpan.Zero, maxTopicSize = 0) | Auto-create the topic on InitAsync if missing |
WithRetry(enabled = true, maxAttempts = 3, initialDelay = 100ms, maxDelay = 10s, backoffMultiplier = 2.0) | Retry failed sends with exponential backoff |
WithBackgroundSending(enabled = true, queueCapacity = 10000, batchSize = 100, flushInterval = 100ms, disposalTimeout = 5s) | Queue and flush sends in the background for higher throughput |
WithEncryptor(encryptor) | Client-side payload encryption. Only valid on a builder-created client — for an external client, set IggyClientConfigurator.MessageEncryptor instead |
SubscribeOnBackgroundError(handler) | Observe background-processing errors (only fires when background sending is enabled) |
SubscribeOnMessageBatchFailed(handler) | Fires when a batch fails after all retries are exhausted (only with background sending — direct sends throw to the caller instead) |
WithLogger(loggerFactory) | Logger factory for diagnostics |
WithRetry parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
enabled | bool | true | Whether retry is enabled |
maxAttempts | int | 3 | Maximum retry attempts |
initialDelay | TimeSpan? | 100ms | Delay before the first retry |
maxDelay | TimeSpan? | 10s | Maximum delay between retries |
backoffMultiplier | double | 2.0 | Exponential backoff multiplier |
WithBackgroundSending parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
enabled | bool | true | Whether background sending is enabled |
queueCapacity | int | 10000 | Max queued send calls (one slot per SendMessagesAsync call, regardless of batch size) |
batchSize | int | 100 | Messages sent per batch |
flushInterval | TimeSpan? | 100ms | Interval at which pending messages are flushed |
disposalTimeout | TimeSpan? | 5s | How long DisposeAsync waits for the background processor to drain |
A batch also flushes once its accumulated payload reaches IggyPublisherConfig.BackgroundMaxBatchBytes (default 256 KB, 0 disables the byte gate) — whichever of batchSize or the byte limit is hit first. This knob has no builder method; set it on the config directly if needed.
Publishing with rented batches
SendMessagesAsync takes messages whose payloads each live in their own byte[]. On hot paths, build a RentedMessageBatch instead — every payload is written into a single buffer rented from the shared array pool — and hand it to SendAsync:
using System.Text.Json;
using Apache.Iggy.Messages;
using var builder = new RentedMessageBatchBuilder(sizeHint: 4096);
foreach (var evt in events)
{
// Serialize straight into the pooled buffer (static lambda + state = no closure allocation)
builder.Add(evt, static (e, writer) =>
{
using var json = new Utf8JsonWriter(writer);
JsonSerializer.Serialize(json, e);
}, Guid.NewGuid());
}
var batch = builder.Build();
await publisher.SendAsync(batch);RentedMessageBatchBuilder:
| Member | Description |
|---|---|
new RentedMessageBatchBuilder(sizeHint = 1024) | Initial capacity hint (bytes) for the rented buffer |
Add(payload, id = null, userHeaders = null) | Copies a ReadOnlySpan<byte> payload into the shared buffer |
Add(state, writePayload, id = null, userHeaders = null) | Writes the payload via an Action<TState, IBufferWriter<byte>> callback directly into the buffer |
Count | Number of payloads added so far |
Build() | Materializes the RentedMessageBatch; the builder must not be used afterwards |
A null message id sends the message with id 0, letting the server assign one.
Ownership rules:
publisher.SendAsync(batch)takes ownership — do not dispose the batch yourself. Its buffer returns to the pool once the batch is sent (immediately, or after the background flush).- For a direct low-level send, pass
batch.Messagestoclient.SendMessagesAsync(...)yourself,awaitit, then dispose the batch (using var batch = builder.Build()). - After disposal,
batch.Messagesand every payload'sSpanthrowObjectDisposedExceptioninstead of reading recycled pool memory. - Dispose the builder only when abandoning it before
Build(); ausingon the builder is safe either way (it becomes a no-op afterBuild()).
Typed publisher
For automatic object serialization, use IggyPublisherBuilder<T> with an ISerializer<T>:
class OrderSerializer : ISerializer<Order>
{
public void Serialize(Order data, IBufferWriter<byte> writer) =>
writer.Write(JsonSerializer.SerializeToUtf8Bytes(data));
}
var publisher = IggyPublisherBuilder<Order>.Create(
client,
Identifier.String("orders-stream"),
Identifier.String("orders-topic"),
new OrderSerializer()
).Build();
await publisher.InitAsync();
await publisher.SendAsync(new List<Order> { /* ... */ });Besides the collection overload, SendAsync also accepts a single item (SendAsync(order, messageId: null, userHeaders: null)) or a collection of (data, messageId, userHeaders) tuples when you need per-message ids or headers.
For JSON you don't need a custom serializer — the built-in SystemTextJsonSerializer<T> (in Apache.Iggy.Publishers) writes System.Text.Json output directly into the send buffer, with optional JsonSerializerOptions:
var publisher = IggyPublisherBuilder<Order>.Create(
client,
Identifier.String("orders-stream"),
Identifier.String("orders-topic"),
new SystemTextJsonSerializer<Order>()
).Build();There is no built-in deserializer counterpart — consumers implement IDeserializer<T> themselves (see Typed consumer).
IggyConsumer
Configure a consumer, initialize it, then iterate. ReceiveAsync returns an IAsyncEnumerable<ReceivedMessage>:
using System.Text;
using Apache.Iggy;
using Apache.Iggy.Consumers;
using Apache.Iggy.Kinds;
var consumer = IggyConsumerBuilder.Create(
client,
Identifier.String("my-stream"),
Identifier.String("my-topic"),
Consumer.New(1)
)
.WithPollingStrategy(PollingStrategy.Next())
.WithBatchSize(10)
.WithAutoCommitMode(AutoCommitMode.Auto)
.Build();
await consumer.InitAsync();
await foreach (var message in consumer.ReceiveAsync())
{
var payload = Encoding.UTF8.GetString(message.Message.Payload);
Console.WriteLine($"Offset {message.CurrentOffset}: {payload}");
}Consumer builder options
| Method | Description |
|---|---|
WithConnection(protocol, address, login, password, receiveBufferSize = 4096, sendBufferSize = 4096, reconnectionSettings = null) | Connection settings — only used when the builder creates its own client |
WithPartitionId(partitionId) | Consume from a specific partition |
WithPollingStrategy(pollingStrategy) | Where to start consuming (default: Offset(0)). An Offset(...) strategy is advanced client-side after each poll; other strategies (e.g. Next()) are sent as-is and rely on server-side offset tracking |
WithBatchSize(batchSize) | Messages fetched per poll (default: 100) |
WithAutoCommitMode(mode) | Offset auto-commit behavior (see below) |
WithConsumerGroup(groupName, createIfNotExists = true, joinGroup = true) | Create and/or join a consumer group |
WithPollingInterval(interval) | Delay between polls to throttle server requests (default: 100ms); TimeSpan.Zero disables throttling |
WithEncryptor(encryptor) | Client-side payload decryption. Only valid on a builder-created client — for an external client, set IggyClientConfigurator.MessageEncryptor instead. Cannot be combined with AutoCommitMode.Auto (the builder throws — a decryption failure would silently skip an already-committed batch) |
SubscribeOnPollingError(handler) | Observe polling errors |
WithLogger(loggerFactory) | Logger factory for diagnostics |
Auto-commit modes (AutoCommitMode):
| Mode | Description |
|---|---|
Auto | Commit the offset while polling |
AfterReceive | Commit after each message is received |
Disabled | Commit manually |
Manual offset control
With AutoCommitMode.Disabled, commit offsets yourself via the consumer's StoreOffsetAsync:
await foreach (var message in consumer.ReceiveAsync())
{
Process(message);
await consumer.StoreOffsetAsync(message.CurrentOffset, message.PartitionId);
}StoreOffsetAsync(offset, partitionId, resetLastPolled = false) stores the offset for a partition; pass resetLastPolled: true to also move the consumer's cached last-polled position so the next poll resumes past the stored offset. DeleteOffsetAsync(partitionId) clears the stored offset.
Consumer groups
Pass a Consumer.Group(...) and let the builder create and join the group for load-balanced consumption:
var consumer = IggyConsumerBuilder.Create(
client,
Identifier.String("my-stream"),
Identifier.String("my-topic"),
Consumer.Group("my-group")
)
.WithConsumerGroup("my-group", createIfNotExists: true, joinGroup: true)
.WithPollingStrategy(PollingStrategy.Next())
.WithAutoCommitMode(AutoCommitMode.AfterReceive)
.Build();
await consumer.InitAsync();
await foreach (var message in consumer.ReceiveAsync())
{
var payload = Encoding.UTF8.GetString(message.Message.Payload);
Console.WriteLine($"Partition {message.PartitionId}: {payload}");
}
await consumer.DisposeAsync();Consuming with rented buffers
ReceiveAsync copies each payload into its own byte[]. ReceiveRentedAsync is the high-level counterpart of the low-level rented poll API: payloads are slices of a pooled buffer shared by all messages from the same poll. Each yielded ReceivedRentedMessage must be disposed — the buffer returns to the pool once the last message of its batch is disposed:
await foreach (var message in consumer.ReceiveRentedAsync())
{
using (message)
{
var text = Encoding.UTF8.GetString(message.Message.Payload.Span);
Console.WriteLine($"Offset {message.CurrentOffset}: {text}");
}
}ReceivedRentedMessage (IDisposable) exposes Message (a RentedMessageResponse), CurrentOffset, PartitionId, Status, and Error. This path performs no deserialization, so Status is always Success. Payload and raw-header memory are only valid until the message is disposed — copy out anything you need to keep (e.g. message.Message.Payload.ToArray()).
Auto-commit modes apply exactly as with ReceiveAsync. Forgetting to dispose a message never corrupts data — the batch buffer simply isn't returned to the pool and is reclaimed by the GC as an ordinary allocation.
Typed consumer
For automatic deserialization, use IggyConsumerBuilder<T> with an IDeserializer<T> and iterate with ReceiveDeserializedAsync. Each ReceivedMessage<T> carries a Status you should check:
class OrderDeserializer : IDeserializer<OrderEvent>
{
public OrderEvent Deserialize(ReadOnlyMemory<byte> data) =>
JsonSerializer.Deserialize<OrderEvent>(data.Span)!;
}
var builder = IggyConsumerBuilder<OrderEvent>.Create(
client,
Identifier.String("orders-stream"),
Identifier.String("orders-topic"),
Consumer.Group("order-processors"),
new OrderDeserializer()
);
builder.WithAutoCommitMode(AutoCommitMode.AfterReceive);
var consumer = builder.Build();
await consumer.InitAsync();
await foreach (var message in consumer.ReceiveDeserializedAsync())
{
if (message.Status == MessageStatus.Success)
{
Console.WriteLine($"Order: {message.Data?.OrderId}");
}
}Note: On the typed builders, the fluent
With*methods are inherited from the untyped base builder and return the base type, while the typedBuild()hides the base one. ChainingCreate(...).WithAutoCommitMode(...).Build()therefore resolves to the baseBuild()and returns an untyped consumer/publisher. Keep the typed builder in a variable (as above) and callBuild()on it — theWith*calls mutate the builder, so their return value can be ignored.