NexJob.Kafka
5.7.0
dotnet add package NexJob.Kafka --version 5.7.0
NuGet\Install-Package NexJob.Kafka -Version 5.7.0
<PackageReference Include="NexJob.Kafka" Version="5.7.0" />
<PackageVersion Include="NexJob.Kafka" Version="5.7.0" />
<PackageReference Include="NexJob.Kafka" />
paket add NexJob.Kafka --version 5.7.0
#r "nuget: NexJob.Kafka, 5.7.0"
#:package NexJob.Kafka@5.7.0
#addin nuget:?package=NexJob.Kafka&version=5.7.0
#tool nuget:?package=NexJob.Kafka&version=5.7.0
NexJob.Kafka
End-to-end Apache Kafka integration for NexJob, unifying Kafka Triggers (Consumers) and a Resilient Outbox Producer in a single package.
Installation
dotnet add package NexJob.Kafka
1. Resilient Outbox Producer
Publishes messages to Kafka topics backed by NexJob's persistent storage, retry with backoff and jitter, dead-letter dispatch, and OpenTelemetry trace propagation.
Registration
using NexJob;
using NexJob.Kafka;
builder.Services.AddNexJobPostgres(builder.Configuration.GetConnectionString("NexJobConnection")!); // or any other storage provider
builder.Services.AddNexJob()
.AddKafkaProducer(options =>
{
options.BootstrapServers = builder.Configuration["KAFKA_BOOTSTRAP_SERVERS"]
?? Environment.GetEnvironmentVariable("KAFKA_BOOTSTRAP_SERVERS")
?? "localhost:9092";
options.Acks = Confluent.Kafka.Acks.All;
options.EnableIdempotence = true;
options.FlushTimeout = TimeSpan.FromSeconds(10);
});
Publishing Messages
Inject IScheduler into your services or controllers:
public class OrderService(IScheduler scheduler)
{
// 1. Strongly typed object (automatically serialized to JSON)
public async Task CreateOrderAsync(OrderCreatedEvent order, CancellationToken ct)
{
await scheduler.EnqueueKafkaAsync(
topic: "orders-topic",
key: order.OrderId.ToString(),
value: order,
cancellationToken: ct);
}
// 2. Raw string / JSON payload
public async Task PublishRawJsonAsync(string orderId, string rawJson, CancellationToken ct)
{
await scheduler.EnqueueKafkaRawAsync(
topic: "orders-topic",
key: orderId,
value: rawJson,
cancellationToken: ct);
}
// 3. Raw binary payload
public async Task PublishBinaryAsync(string key, byte[] bytes, CancellationToken ct)
{
await scheduler.EnqueueKafkaRawAsync(
topic: "events-topic",
key: key,
value: bytes,
cancellationToken: ct);
}
}
2. Kafka Trigger (Consumer)
Consumes incoming messages from Kafka topics and automatically enqueues them as background jobs.
Registration
There are two ways to register which job is triggered when a message arrives:
Option A: Strongly-Typed Consumer (Recommended)
Bind a topic directly to a job handler class (IJob<string>). The job is automatically registered in DI as Transient:
builder.Services.AddNexJob()
.AddKafkaTrigger<ProcessOrderJob>(options =>
{
options.BootstrapServers = builder.Configuration["KAFKA_BOOTSTRAP_SERVERS"] ?? "localhost:9092";
options.Topic = "incoming-orders";
options.GroupId = "nexjob-consumer-group";
options.TargetQueue = "orders";
});
Option B: Dynamic Message Header (nexjob.job_type)
If multiple job types share the same topic, omit the generic argument. The incoming Kafka message must include the nexjob.job_type header (or have options.JobType set as a default):
builder.Services.AddNexJob()
.AddKafkaTrigger(options =>
{
options.BootstrapServers = builder.Configuration["KAFKA_BOOTSTRAP_SERVERS"] ?? "localhost:9092";
options.Topic = "incoming-orders";
options.GroupId = "nexjob-consumer-group";
options.TargetQueue = "orders";
// options.JobType = typeof(DefaultOrderJob).AssemblyQualifiedName;
});
Inbound Message Contract
For dynamic triggers (Option B), consumed messages can carry the following headers:
nexjob.job_type: Assembly-qualified name of theIJob<string>to execute. When it is absent,options.JobTypeis used; a message with neither can never become a job and is handled as a permanent failure.traceparent: W3C distributed trace header (optional).
The message value is passed as the string input to the resolved job. The job idempotency key is the record position (kafka:{topic}:{partition}:{offset}), so a redelivered record never creates a second job while two records that share a message key each produce a job.
Enqueue failures are classified: a transient failure (storage or network error) retries the same record in place with a 1 s to 30 s backoff and commits only after success; a permanent failure goes to DeadLetterTopic and is committed (without a topic it is logged at Error and committed, so it cannot block the partition).
3. Configuration & 12-Factor App (Docker / Kubernetes)
NexJob.Kafka is configured through the option delegates shown above; it does not read appsettings.json sections or environment variables by itself. In containerized environments, read the variables in your own code and assign them to the options:
# Set environment variables in Docker / Kubernetes
export KAFKA_BOOTSTRAP_SERVERS="kafka-broker.prod:9092"
export KAFKA_TOPIC="orders"
Read them in C#:
builder.Services.AddNexJob()
.AddKafkaProducer(options =>
{
options.BootstrapServers = Environment.GetEnvironmentVariable("KAFKA_BOOTSTRAP_SERVERS")
?? throw new InvalidOperationException("KAFKA_BOOTSTRAP_SERVERS is not set.");
});
4. Configuration Options
Producer Options (KafkaProducerOptions)
| Option | Description | Default |
|---|---|---|
BootstrapServers |
Comma-separated list of Kafka broker endpoints (required) | "" |
Acks |
Acknowledgment guarantee (Acks.All, Acks.Leader, Acks.None) |
Acks.All |
EnableIdempotence |
Controls producer idempotence on the broker | true |
FlushTimeout |
Timeout for flushing messages on application shutdown | 10 seconds |
Queue |
NexJob queue name used for publishing jobs | "kafka-producer" |
DefaultPriority |
Default execution priority for publishing jobs | JobPriority.Normal |
ConfigureProducer |
Delegate (Action<ProducerConfig>) to customize SASL/SSL credentials, TLS certificates, and timeouts |
null |
Trigger Options (KafkaTriggerOptions)
| Option | Description | Default |
|---|---|---|
BootstrapServers |
Comma-separated list of Kafka broker endpoints (required) | "" |
Topic |
Kafka topic to consume messages from | "" |
GroupId |
Kafka consumer group identifier | "" |
TargetQueue |
Target NexJob queue name for enqueued jobs | "default" |
JobPriority |
Priority of the enqueued jobs | JobPriority.Normal |
DeadLetterTopic |
Topic that receives records that can never be enqueued (permanent failures) | null |
JobType |
Assembly-qualified job type used when a message has no nexjob.job_type header |
null |
ConsumeTimeout |
Polling timeout for IConsumer.Consume |
1 second |
ConfigureConsumer |
Delegate (Action<ConsumerConfig>) to customize SASL/SSL credentials, TLS certificates, and timeouts |
null |
5. Architectural Guarantees
- Transactional Outbox / Resilient Publish: Messages are first durably persisted to NexJob storage. If the Kafka broker is down, messages wait safely in storage and are retried automatically.
- Dead-Letter Handling: Permanent publishing failures trigger NexJob's dead-letter pipeline (
IDeadLetterHandler) and surface in the dashboard. - Trace Propagation: Injects W3C
traceparentheaders into outgoing messages and extracts them on consumer triggers. - Graceful Shutdown: Unflushed in-flight messages are flushed before the application process exits.
6. Sagas & Event-Driven Choreographies
Need distributed state machines or sagas with compensating transactions over Kafka? See qKafka — the companion event-driven framework that pairs with NexJob.
Forwarding failed jobs
A job created from a consumed message that exhausts its retries stays only inside NexJob (Failed in the dashboard), because the broker message was acknowledged when it became a job. To forward a copy of the original message to a Kafka topic, see Forwarding a dead-lettered job.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net8.0
- Confluent.Kafka (>= 2.14.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 8.0.1)
- Microsoft.Extensions.Logging.Abstractions (>= 8.0.3)
- Microsoft.Extensions.Options (>= 8.0.2)
- Microsoft.Extensions.Options.DataAnnotations (>= 8.0.0)
- NexJob (>= 5.7.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.