Table of Contents

Kafka - File Streaming

This sample demonstrates how to deal with raw binary contents and large messages, to transfer some files through Kafka.

See also: Serializing the Produced Messages, Deserializing the Consumed Messages, Producing Chunked Messages, Consuming Chunked Messages

Producer

The producer exposes two REST API that receive the path of a local file to be streamed. The second API uses a custom BinaryMessage to forward further metadata (the file name in this example).

using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;
using Silverback.Configuration;
using Silverback.Messaging.Configuration;

namespace Silverback.Samples.Kafka.BinaryFileStreaming.Producer;

public class Startup
{
    public void ConfigureServices(IServiceCollection services)
    {
        // Enable Silverback
        services.AddSilverback()

            // Use Apache Kafka as message broker
            .WithConnectionToMessageBroker(options => options.AddKafka())

            // Delegate the broker clients configuration to a separate class
            .AddBrokerClientsConfigurator<BrokerClientsConfigurator>();

        // Add API controllers and SwaggerGen
        services.AddControllers();
        services.AddSwaggerGen();
    }

    public void Configure(IApplicationBuilder app)
    {
        // Enable middlewares to serve generated Swagger JSON and UI
        app.UseSwagger().UseSwaggerUI(
            uiOptions =>
            {
                uiOptions.SwaggerEndpoint(
                    "/swagger/v1/swagger.json",
                    $"{GetType().Assembly.FullName} API");
            });

        // Enable routing and endpoints for controllers
        app.UseRouting();
        app.UseEndpoints(endpoints => { endpoints.MapControllers(); });
    }
}

Full source code: https://github.com/BEagle1984/silverback/tree/master/samples/Kafka/BinaryFileStreaming.Producer

Consumer

The consumer simply streams the file to a temporary folder in the local file system.

using Microsoft.Extensions.DependencyInjection;
using Silverback.Configuration;
using Silverback.Messaging.Configuration;
using Silverback.Samples.Kafka.BinaryFileStreaming.Consumer.Subscribers;

namespace Silverback.Samples.Kafka.BinaryFileStreaming.Consumer;

public class Startup
{
    public void ConfigureServices(IServiceCollection services)
    {
        // Enable Silverback
        services.AddSilverback()

            // Use Apache Kafka as message broker
            .WithConnectionToMessageBroker(options => options.AddKafka())

            // Delegate the broker clients configuration to a separate class
            .AddBrokerClientsConfigurator<BrokerClientsConfigurator>()

            // Register the subscribers
            .AddSingletonSubscriber<BinaryFileSubscriber>();
    }

    public void Configure()
    {
    }
}

Full source code: https://github.com/BEagle1984/silverback/tree/master/samples/Kafka/BinaryFileStreaming.Consumer