// Objects on the wire: single-object messages, schema-registry framing, and streams of objects without a container.
using System;
using System.Collections.Generic;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using AvroSharp.Generic;
using AvroSharp.Messages;
using AvroSharp.Schemas;
using AvroSharp.Streams;
var v1 = (RecordSchema)AvroSchema.Parse("""{"type":"record","name":"Reading","namespace":"iot","fields":[{"name":"sensor","type":"string"},{"name":"value","type":"float"}]}""");
var v2 = (RecordSchema)AvroSchema.Parse("""{"type":"record","name":"Reading","namespace":"iot","fields":[{"name":"sensor","type":"string"},{"name":"value","type":"double"},{"name":"unit","type":"string","default":"C"}]}""");
var reading = new GenericRecord(v1) { ["sensor"] = "t-1", ["value"] = 21.5f };
var writer = GenericDatumWriter.Create(v1);
// Single-object encoding: C3 01, the writer schema's CRC-64 fingerprint, then the data. The reader finds the schema
// by its fingerprint and resolves every message to the schema it asks for.
byte[] message = AvroMessage.ToArray(reading, writer);
var messages = AvroMessageReader.CreateGeneric(new AvroSchemaStore(v1, v2), readerSchema: v2);
var fromMessage = messages.Read(message).AsRecord();
Console.WriteLine($"Single-object message: {message.Length} bytes -> value {fromMessage["value"].AsDouble()}, unit {fromMessage["unit"].AsString()}");
// Schema-registry framing (here Confluent's: 0x00 and a 4-byte schema ID). The resolver maps IDs to schemas; a real
// one would call the registry and cache the answers.
var registry = new Messaging.InMemoryRegistry { [1] = v1, [2] = v2 };
byte[] framed = AvroRegistryMessage.ToArray(AvroRegistryFraming.Confluent, AvroSchemaId.FromNumber(1), reading, writer);
var framedReader = AvroRegistryMessageReader.CreateGeneric(AvroRegistryFraming.Confluent, registry);
var fromRegistry = (await framedReader.ReadAsync(framed)).AsRecord();
Console.WriteLine($"Confluent-framed message: {framed.Length} bytes, schema ID 1 -> sensor {fromRegistry["sensor"].AsString()}");
// A stream of objects with no container or framing, for a socket or a file of concatenated objects. The schema is
// not in the stream: both sides must know it.
using var stream = new MemoryStream();
await using (var streamWriter = AvroStreamWriter.CreateGeneric(stream, v1, new AvroStreamOptions { LeaveOpen = true }))
{
for (var i = 0; i < 100; i++)
{
await streamWriter.WriteAsync(new GenericRecord(v1) { ["sensor"] = $"t-{i % 4}", ["value"] = i / 2f });
}
}
stream.Position = 0;
var total = 0.0;
var count = 0;
await using (var streamReader = AvroStreamReader.OpenGeneric(stream, writerSchema: v1, readerSchema: v2, new AvroStreamOptions { LeaveOpen = true }))
{
await foreach (var value in streamReader.ReadAllAsync())
{
total += value.AsRecord()["value"].AsDouble();
count++;
}
}
Console.WriteLine($"Stream of objects: {stream.Length} bytes, {count} readings, sum {total}");
var ok = Math.Abs(fromMessage["value"].AsDouble() - 21.5) < 1e-9 && string.Equals(fromMessage["unit"].AsString(), "C", StringComparison.Ordinal)
&& fromRegistry.Equals(reading)
&& count == 100 && Math.Abs(total - 2475) < 1e-9;
Console.WriteLine(ok ? "OK" : "FAILED");
return ok ? 0 : 1;
namespace Messaging
{
/// Schema IDs to schemas, in memory. A registry client would fetch unknown IDs and cache them.
internal sealed class InMemoryRegistry : Dictionary, IAvroSchemaIdResolver
{
public AvroSchema? GetSchema(AvroSchemaId id) => TryGetValue(id.Number, out var schema) ? schema : null;
public ValueTask GetSchemaAsync(AvroSchemaId id, CancellationToken cancellationToken = default) => new(GetSchema(id));
}
}