|
| 1 | +// Licensed to the .NET Foundation under one or more agreements. |
| 2 | +// The .NET Foundation licenses this file to you under the MIT license. |
| 3 | + |
| 4 | +using Aspire.Hosting.ApplicationModel; |
| 5 | +using Aspire.Hosting.Publishing; |
| 6 | + |
| 7 | +namespace Aspire.Hosting; |
| 8 | + |
| 9 | +public static class KafkaBuilderExtensions |
| 10 | +{ |
| 11 | + private const int KafkaBrokerPort = 9092; |
| 12 | + /// <summary> |
| 13 | + /// Adds a Kafka broker container to the application. |
| 14 | + /// </summary> |
| 15 | + /// <param name="builder">The <see cref="IDistributedApplicationBuilder"/>.</param> |
| 16 | + /// <param name="name">The name of the resource. This name will be used as the connection string name when referenced in a dependency.</param> |
| 17 | + /// <param name="port">The host port of Kafka broker.</param> |
| 18 | + /// <returns>A reference to the <see cref="IResourceBuilder{KafkaContainerResource}"/></returns> |
| 19 | + public static IResourceBuilder<KafkaContainerResource> AddKafkaContainer(this IDistributedApplicationBuilder builder, string name, int? port = null) |
| 20 | + { |
| 21 | + var kafka = new KafkaContainerResource(name); |
| 22 | + return builder.AddResource(kafka) |
| 23 | + .WithEndpoint(hostPort: port, containerPort: KafkaBrokerPort) |
| 24 | + .WithAnnotation(new ContainerImageAnnotation { Image = "confluentinc/confluent-local", Tag = "latest" }) |
| 25 | + .WithManifestPublishingCallback(context => WriteKafkaContainerToManifest(context, kafka)) |
| 26 | + .WithEnvironment(context => ConfigureKafkaContainer(context, kafka)); |
| 27 | + |
| 28 | + static void WriteKafkaContainerToManifest(ManifestPublishingContext context, KafkaContainerResource resource) |
| 29 | + { |
| 30 | + context.WriteContainer(resource); |
| 31 | + context.Writer.WriteString("connectionString", $"{{{resource.Name}.bindings.tcp.host}}:{{{resource.Name}.bindings.tcp.port}}"); |
| 32 | + } |
| 33 | + } |
| 34 | + |
| 35 | + /// <summary> |
| 36 | + /// Adds a Kafka resource to the application. A container is used for local development. |
| 37 | + /// </summary> |
| 38 | + /// <param name="builder">The <see cref="IDistributedApplicationBuilder"/>.</param> |
| 39 | + /// <param name="name">The name of the resource. This name will be used as the connection string name when referenced in a dependency</param> |
| 40 | + /// <returns>A reference to the <see cref="IResourceBuilder{KafkaServerResource}"/>.</returns> |
| 41 | + public static IResourceBuilder<KafkaServerResource> AddKafka(this IDistributedApplicationBuilder builder, string name) |
| 42 | + { |
| 43 | + var kafka = new KafkaServerResource(name); |
| 44 | + return builder.AddResource(kafka) |
| 45 | + .WithEndpoint(containerPort: KafkaBrokerPort) |
| 46 | + .WithAnnotation(new ContainerImageAnnotation{ Image = "confluentinc/confluent-local", Tag = "latest" }) |
| 47 | + .WithManifestPublishingCallback(WriteKafkaServerToManifest) |
| 48 | + .WithEnvironment(context => ConfigureKafkaContainer(context, kafka)); |
| 49 | + |
| 50 | + static void WriteKafkaServerToManifest(ManifestPublishingContext context) |
| 51 | + { |
| 52 | + context.Writer.WriteString("type", "kafka.server.v0"); |
| 53 | + } |
| 54 | + } |
| 55 | + |
| 56 | + private static void ConfigureKafkaContainer(EnvironmentCallbackContext context, IResource resource) |
| 57 | + { |
| 58 | + // confluentinc/confluent-local is a docker image that contains a Kafka broker started with KRaft to avoid pulling a separate image for ZooKeeper. |
| 59 | + // See https://github.com/confluentinc/kafka-images/blob/master/local/README.md. |
| 60 | + // When not explicitly set default configuration is applied. |
| 61 | + // See https://github.com/confluentinc/kafka-images/blob/master/local/include/etc/confluent/docker/configureDefaults for more details. |
| 62 | + |
| 63 | + var hostPort = context.PublisherName == "manifest" |
| 64 | + ? KafkaBrokerPort |
| 65 | + : GetResourcePort(resource); |
| 66 | + context.EnvironmentVariables.Add("KAFKA_ADVERTISED_LISTENERS", |
| 67 | + $"PLAINTEXT://localhost:29092,PLAINTEXT_HOST://localhost:{hostPort}"); |
| 68 | + |
| 69 | + static int GetResourcePort(IResource resource) |
| 70 | + { |
| 71 | + if (!resource.TryGetAllocatedEndPoints(out var allocatedEndpoints)) |
| 72 | + { |
| 73 | + throw new DistributedApplicationException( |
| 74 | + $"Kafka resource \"{resource.Name}\" does not have endpoint annotation."); |
| 75 | + } |
| 76 | + |
| 77 | + return allocatedEndpoints.Single().Port; |
| 78 | + } |
| 79 | + } |
| 80 | +} |
0 commit comments