Browse Source

upgrade Confluent.Kafka to 1.4.3

master
Savorboard 4 years ago
parent
commit
6f61c4c3f4
2 changed files with 5 additions and 5 deletions
  1. +1
    -1
      src/DotNetCore.CAP.Kafka/DotNetCore.CAP.Kafka.csproj
  2. +4
    -4
      src/DotNetCore.CAP.Kafka/KafkaConsumerClient.cs

+ 1
- 1
src/DotNetCore.CAP.Kafka/DotNetCore.CAP.Kafka.csproj View File

@@ -13,7 +13,7 @@
</PropertyGroup>

<ItemGroup>
<PackageReference Include="Confluent.Kafka" Version="1.3.0" />
<PackageReference Include="Confluent.Kafka" Version="1.4.3" />
</ItemGroup>

<ItemGroup>


+ 4
- 4
src/DotNetCore.CAP.Kafka/KafkaConsumerClient.cs View File

@@ -52,10 +52,10 @@ namespace DotNetCore.CAP.Kafka
{
var consumerResult = _consumerClient.Consume(cancellationToken);

if (consumerResult.IsPartitionEOF || consumerResult.Value == null) continue;
if (consumerResult.IsPartitionEOF || consumerResult.Message.Value == null) continue;

var headers = new Dictionary<string, string>(consumerResult.Headers.Count);
foreach (var header in consumerResult.Headers)
var headers = new Dictionary<string, string>(consumerResult.Message.Headers.Count);
foreach (var header in consumerResult.Message.Headers)
{
var val = header.GetValueBytes();
headers.Add(header.Key, val != null ? Encoding.UTF8.GetString(val) : null);
@@ -71,7 +71,7 @@ namespace DotNetCore.CAP.Kafka
}
}

var message = new TransportMessage(headers, consumerResult.Value);
var message = new TransportMessage(headers, consumerResult.Message.Value);

OnMessageReceived?.Invoke(consumerResult, message);
}


Loading…
Cancel
Save