KafkaService 1.0.2
The owner has unlisted this package.
This could mean that the package is deprecated, has security vulnerabilities or shouldn't be used anymore.
dotnet add package KafkaService --version 1.0.2
NuGet\Install-Package KafkaService -Version 1.0.2
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="KafkaService" Version="1.0.2" />
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="KafkaService" Version="1.0.2" />
<PackageReference Include="KafkaService" />
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add KafkaService --version 1.0.2
The NuGet Team does not provide support for this client. Please contact its maintainers for support.
#r "nuget: KafkaService, 1.0.2"
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package KafkaService@1.0.2
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=KafkaService&version=1.0.2
#tool nuget:?package=KafkaService&version=1.0.2
The NuGet Team does not provide support for this client. Please contact its maintainers for support.
Kafka Service
Consumer side settings
Program.cs file
builder.Services.AddSingleton<WebSocketConnectionManager>();
builder.Services.AddSingleton<IKafaProducer, KafkaProducer>();
builder.Services.AddSingleton<IKafkaMessageHandler, ConsumerKafkaHandler>();
builder.Services.AddSingleton<IKafkaConsumer, KafkaConsumer>();
builder.Services.AddHostedService<KafkaConsumer>();
appsettings.join
"Kafka": {
"BootstrapServers": "localhost:9092",
"Topics": [ "PostRecord"]
},
KafkaHandler.cs file
public class ConsumerKafkaHandler : IKafkaMessageHandler
{
private readonly IKafaProducer kafaProducer;
public ConsumerKafkaHandler(IKafaProducer kafaProducer)
{
this.kafaProducer = kafaProducer;
}
public async Task<Task> HandleMessageAsync(string topic, string message)
{
await Task.Delay(0);
Console.WriteLine($"Consumer: Received message on topic {topic}: {message}");
MyApplicationMethod(message);
switch (topic)
{
case "PostRecord":
var res = kafaProducer.ProduceAsync("PostRecord_ACK", "PostRecord_ACK", message);
break;
default:
Console.WriteLine($"Unknown topic: {topic}");
break;
}
return Task.CompletedTask;
}
private void MyApplicationMethod(string msg)
{
//Console.WriteLine($"Handling message in app: {msg}");
}
}
Producer side settings
Program.cs file
builder.Services.AddSingleton<WebSocketConnectionManager>();
builder.Services.AddSingleton<IKafkaMessageHandler, ProducerKafkaHandler>();
builder.Services.AddSingleton<IKafkaConsumer, KafkaConsumer>();
builder.Services.AddSingleton<IKafaProducer,KafkaProducer>();
builder.Services.AddHostedService<KafkaConsumer>();
app.UseWebSockets();
appsettings.join
"Kafka": {
"BootstrapServers": "localhost:9092",
"Topics": [ "PostRecord_ACK" ]
},
KafkaHandler.cs file
public class ProducerKafkaHandler : IKafkaMessageHandler
{
private readonly WebSocketConnectionManager _connectionManager;
public MyKafkaHandler(WebSocketConnectionManager connectionManager)
{
_connectionManager = connectionManager;
}
public async Task<Task> HandleMessageAsync(string topic, string message)
{
// Send ack via WebSocket to client
await _connectionManager.SendMessageAsync("PostRecord", $"Saved {topic} : {message}");
//_connectionManager.RemoveSocket("PostRecord");
// Call your application's logic here
Console.WriteLine($"Producer: Received message on topic {topic}: {message}");
MyApplicationMethod(message);
return Task.CompletedTask;
}
private void MyApplicationMethod(string msg)
{
}
}
Producer Controller
[Route("api/[controller]")]
[ApiController]
public class ProducerController : ControllerBase
{
private readonly IKafaProducer kafaProducer;
public ProducerController(IKafaProducer kafaProducer)
{
this.kafaProducer = kafaProducer;
}
[HttpPost]
public Task<IActionResult> Post([FromBody] PostRequest request)
{
var message = JsonSerializer.Serialize(request);
var res = kafaProducer.ProduceAsync("PostRecord", "PostRecord", message);
return Task.FromResult<IActionResult>(Ok("Record Updated Successfully..."));
}
}
WebSocketController
[ApiController]
[Route("ws")]
public class WebSocketController : ControllerBase
{
private readonly WebSocketConnectionManager _connectionManager;
public WebSocketController(WebSocketConnectionManager connectionManager)
{
_connectionManager = connectionManager;
}
[HttpGet()]
public async Task Get()
{
string clientId = "PostRecord";
if (HttpContext.WebSockets.IsWebSocketRequest)
{
var socket = await HttpContext.WebSockets.AcceptWebSocketAsync();
_connectionManager.AddSocket(clientId, socket);
var buffer = new byte[1024 * 4];
while (socket.State == WebSocketState.Open)
{
var result = await socket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None);
if (result.MessageType == WebSocketMessageType.Close)
{
_connectionManager.RemoveSocket(clientId);
await socket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closed", CancellationToken.None);
}
}
}
else
{
HttpContext.Response.StatusCode = 400;
}
}
}
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net9.0 is compatible. 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. |
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.
-
net9.0
- Confluent.Kafka (>= 2.10.0)
- Microsoft.Extensions.Hosting.Abstractions (>= 9.0.5)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|
Initial Note