High-performance .NET SDK for streaming data ingestion into Databricks Delta tables using the Zerobus service. Built on the same Rust core as the Go SDK, exposed via P/Invoke (C FFI bindings).
- .NET 8 or .NET 10
- Rust toolchain (for building the native
zerobus_ffilibrary from source)
using Databricks.Zerobus;
// 1. Create SDK instance.
using var sdk = ZerobusSdk.CreateBuilder()
.Endpoint("https://your-shard.zerobus.databricks.com")
.UnityCatalogUrl("https://your-workspace.databricks.com")
.Build();
// 2. Configure stream options.
var options = StreamConfigurationOptions.Default with
{
MaxInflightRequests = 50_000,
};
// 3. Create stream.
using var stream = sdk.CreateJsonStream(
"catalog.schema.table",
clientId,
clientSecret,
options);
// 4. Ingest records.
long offset = stream.IngestRecord("""{"id": 1, "message": "Hello"}""");
// 5. Wait for acknowledgment.
stream.WaitForOffset(offset);dotnet add package Databricks.Zerobuscd dotnet
dotnet buildThe build automatically invokes build_native.sh to compile the Rust FFI shared library and place it in the correct runtimes/<RID>/native/ directory. You need cargo on your PATH (or in ~/.cargo/bin/).
To skip the automatic native build (e.g. when the library is pre-built):
dotnet build -p:SkipNativeBuild=trueThe main entry point. Manages the connection to Zerobus and Unity Catalog.
using var sdk = ZerobusSdk.CreateBuilder()
.Endpoint(zerobusEndpoint)
.UnityCatalogUrl(unityCatalogUrl)
.Build();Creates a JSON-only stream with OAuth 2.0 client credentials authentication.
using var stream = sdk.CreateJsonStream(
"catalog.schema.table",
clientId,
clientSecret,
options); // optional, defaults if nullCreates a protobuf-only stream with OAuth 2.0 client credentials authentication.
using var stream = sdk.CreateProtoStream(
"catalog.schema.table",
descriptorProto,
clientId,
clientSecret,
options); // optional, defaults if nullCreates the legacy untyped stream with OAuth 2.0 client credentials authentication.
Use this only when you intentionally need access to both JSON and protobuf overloads.
The SDK now validates the RecordType/DescriptorProto combination before calling Rust.
using var stream = sdk.CreateStream(
new TableProperties("catalog.schema.table"),
clientId,
clientSecret,
options); // optional, defaults if nullCreates a JSON-only stream with custom authentication headers.
using var stream = sdk.CreateJsonStreamWithHeadersProvider(
"catalog.schema.table",
new MyHeadersProvider(),
options); // optionalCreates a protobuf-only stream with custom authentication headers.
using var stream = sdk.CreateProtoStreamWithHeadersProvider(
"catalog.schema.table",
descriptorProto,
new MyHeadersProvider(),
options); // optionalCreates a stream with custom authentication headers.
using var stream = sdk.CreateStreamWithHeadersProvider(
new TableProperties("catalog.schema.table"),
new MyHeadersProvider(),
options); // optionalRecreates a stream for recovery after the original stream has failed or closed.
Important: RecreateStream(stream) transfers ownership from stream to the returned
stream. The input stream is disposed during recreation and must not be used afterward.
A later Dispose() on the original wrapper (for example at the end of a using scope)
is a no-op.
Typed stream wrappers expose only the matching ingest overloads while preserving the
same lifecycle APIs (Flush, WaitForOffset, Close, Dispose, GetUnackedRecords).
stream.IngestRecord("""{"field": "value"}""");
stream.Flush();byte[] protoBytes = myMessage.ToByteArray();
stream.IngestRecord(protoBytes);
stream.Flush();The untyped stream remains available for advanced callers. Thread-safe.
Ingests a single record and returns its offset immediately (acknowledgment happens in background).
If you use this untyped API, JSON streams must set RecordType.Json and use the string overloads.
Proto streams must provide DescriptorProto and use the byte-oriented overloads.
// JSON
long offset = stream.IngestRecord("""{"field": "value"}""");
// Protobuf
byte[] protoBytes = myMessage.ToByteArray();
long offset = stream.IngestRecord(protoBytes);
stream.WaitForOffset(offset);Ingests a batch of records and returns one offset for the whole batch.
string[] records = [
"""{"device": "sensor-001", "temp": 20}""",
"""{"device": "sensor-002", "temp": 21}""",
];
long lastOffset = stream.IngestRecords(records);
stream.WaitForOffset(lastOffset);Blocks until a specific offset is acknowledged by the server.
stream.WaitForOffset(offset);Blocks until all pending records are acknowledged.
stream.Flush();Retrieves unacknowledged records after stream failure (call after close/failure only).
ReadOnlyMemory<byte>[] unacked = stream.GetUnackedRecords();
foreach (var payload in unacked)
{
// Decode as UTF-8 if you know this stream ingests JSON.
Console.WriteLine($"record bytes: {payload.Length}");
}Close() gracefully closes the stream (flushes first) but keeps the stream
readable for recovery (GetUnackedRecords, RecreateStream).
If you call RecreateStream(stream), the original stream object is disposed and
no longer usable.
Dispose() frees native resources and should be called when recovery work is done.
using calls Dispose() automatically.
stream.Close();
// or let `using` / `await using` handle cleanupInterface for custom authentication.
public class CustomHeadersProvider : IHeadersProvider
{
public IDictionary<string, string> GetHeaders()
{
return new Dictionary<string, string>
{
["authorization"] = "Bearer " + GetToken(),
["x-databricks-zerobus-table-name"] = "catalog.schema.table",
};
}
}Use C# record with expressions to customise:
var options = StreamConfigurationOptions.Default with
{
MaxInflightRequests = 50_000,
RecoveryRetries = 10,
};| Property | Default | Description |
|---|---|---|
MaxInflightRequests |
1,000,000 | Backpressure control |
Recovery |
true |
Auto-recovery on failures |
RecoveryTimeoutMs |
15,000 | Timeout per recovery attempt |
RecoveryBackoffMs |
2,000 | Delay between retries |
RecoveryRetries |
4 | Max recovery attempts |
ServerLackOfAckTimeoutMs |
60,000 | Server ack timeout |
FlushTimeoutMs |
300,000 | Flush timeout (5 min) |
RecordType |
Proto |
Proto / Json / Unspecified |
StreamPausedMaxWaitTimeMs |
null |
Graceful close wait time |
Typed factories set RecordType automatically. You only need to set it manually when using
the untyped CreateStream or CreateStreamWithHeadersProvider APIs.
Errors throw ZerobusException with an IsRetryable property:
try
{
long offset = stream.IngestRecord(data);
}
catch (ZerobusException ex) when (ex.IsRetryable)
{
// Transient error — SDK auto-recovers when Recovery is enabled.
Console.WriteLine($"Retryable error: {ex.RawMessage}");
}
catch (ZerobusException ex)
{
// Fatal error — manual intervention needed.
Console.WriteLine($"Fatal error: {ex.RawMessage}");
}The native zerobus_ffi shared library is built automatically when you run dotnet build. The MSBuild target invokes build_native.sh, which:
- Detects your OS and architecture
- Runs
cargo build --releasein thezerobus-fficrate - Copies the shared library (
.dylib/.so/.dll) tosrc/Zerobus/runtimes/<RID>/native/ - Skips the rebuild if the library is already up to date
You can also run the script directly:
cd dotnet
./build_native.sh # Build for current platform
./build_native.sh --force # Force rebuildThe native library is placed in the standard .NET runtime identifier layout:
| Platform | Path |
|---|---|
| Linux x64 | runtimes/linux-x64/native/libzerobus_ffi.so |
| Linux arm64 | runtimes/linux-arm64/native/libzerobus_ffi.so |
| macOS x64 | runtimes/osx-x64/native/libzerobus_ffi.dylib |
| macOS arm64 | runtimes/osx-arm64/native/libzerobus_ffi.dylib |
| Windows x64 | runtimes/win-x64/native/zerobus_ffi.dll |
Unit tests are isolated and do not require the native library:
dotnet test tests/Zerobus.TestsIntegration tests spin up a mock gRPC server (per test) and exercise the full SDK through the native FFI layer. They require the Rust toolchain to build the native library:
dotnet test tests/Zerobus.IntegrationTestsThe integration tests cover:
| Test | Scenario |
|---|---|
SuccessfulStreamCreation |
Stream creation succeeds |
TimeoutedStreamCreation |
Timeout when server responds slowly |
NonRetriableErrorDuringStreamCreation |
Non-retriable error (e.g. Unauthenticated) |
RetriableErrorWithoutRecoveryDuringStreamCreation |
Retriable error with recovery disabled |
GracefulClose |
Ingest record then close gracefully |
IdempotentClose |
Multiple Close() calls succeed |
IngestAfterClose |
Ingest after close throws |
IngestSingleRecord |
Single record ingest and ack |
IngestMultipleRecords |
Multiple sequential records with ack |
IngestBatchRecords |
Batch ingest of 5 records |
IngestRecordsAfterClose |
Batch ingest after close throws |
Each test gets its own mock gRPC server on a unique port, so all tests run in parallel.
dotnet testdotnet/
├── Zerobus.slnx # Solution file
├── Directory.Build.props # Shared build settings
├── build_native.sh # Rust FFI build script
├── README.md
├── src/
│ └── Zerobus/ # Main SDK library
│ ├── Zerobus.csproj
│ ├── ZerobusSdk.cs # SDK entry point (IDisposable)
│ ├── ZerobusStream.cs # Stream for record ingestion (IDisposable)
│ ├── ZerobusException.cs # Error type with IsRetryable
│ ├── IHeadersProvider.cs # Custom auth interface
│ ├── RecordType.cs # Proto / Json / Unspecified enum
│ ├── StreamConfigurationOptions.cs # Config record with defaults
│ ├── TableProperties.cs # Table name + optional descriptor
│ ├── Properties/
│ │ └── AssemblyInfo.cs
│ └── Native/ # P/Invoke layer (internal)
│ ├── NativeBindings.cs # Raw DllImport declarations
│ ├── NativeInterop.cs # Safe wrappers + marshalling
│ └── HeadersProviderBridge.cs # Managed→native callback bridge
├── tests/
│ ├── Zerobus.Tests/ # Unit tests (NUnit)
│ └── Zerobus.IntegrationTests/ # Integration tests (NUnit + gRPC mock)
│ ├── Zerobus.IntegrationTests.csproj
│ ├── IntegrationTests.cs # 11 integration tests
│ ├── MockZerobusServer.cs # Mock gRPC server
│ ├── TestHelpers.cs # Fixtures, response builders, interceptor
│ └── Protos/
│ └── zerobus_service.proto # Proto definition for gRPC stubs
└── examples/
├── JsonSingle/ # Single JSON record ingestion
├── JsonBatch/ # Batch JSON record ingestion
├── ProtoSingle/ # Single protobuf record ingestion
.NET SDK (Databricks.Zerobus)
↓ P/Invoke
Rust FFI (zerobus-ffi / libzerobus_ffi)
↓
Rust Core (databricks-zerobus-ingest-sdk)
↓ gRPC
Zerobus Service