|
@@ -66,16 +66,18 @@ namespace DotNetCore.CAP.Kafka |
|
|
|
|
|
|
|
|
private void InitKafkaClient() |
|
|
private void InitKafkaClient() |
|
|
{ |
|
|
{ |
|
|
_kafkaOptions.MainConfig["group.id"] = _groupId; |
|
|
|
|
|
|
|
|
lock (_kafkaOptions) |
|
|
|
|
|
{ |
|
|
|
|
|
_kafkaOptions.MainConfig["group.id"] = _groupId; |
|
|
|
|
|
|
|
|
var config = _kafkaOptions.AsKafkaConfig(); |
|
|
|
|
|
_consumerClient = new Consumer<Null, string>(config, null, StringDeserializer); |
|
|
|
|
|
_consumerClient.OnConsumeError += ConsumerClient_OnConsumeError; |
|
|
|
|
|
_consumerClient.OnMessage += ConsumerClient_OnMessage; |
|
|
|
|
|
_consumerClient.OnError += ConsumerClient_OnError; |
|
|
|
|
|
|
|
|
var config = _kafkaOptions.AsKafkaConfig(); |
|
|
|
|
|
_consumerClient = new Consumer<Null, string>(config, null, StringDeserializer); |
|
|
|
|
|
_consumerClient.OnConsumeError += ConsumerClient_OnConsumeError; |
|
|
|
|
|
_consumerClient.OnMessage += ConsumerClient_OnMessage; |
|
|
|
|
|
_consumerClient.OnError += ConsumerClient_OnError; |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private void ConsumerClient_OnConsumeError(object sender, Message e) |
|
|
private void ConsumerClient_OnConsumeError(object sender, Message e) |
|
|
{ |
|
|
{ |
|
|
var message = e.Deserialize<Null, string>(null, StringDeserializer); |
|
|
var message = e.Deserialize<Null, string>(null, StringDeserializer); |
|
|