-
-
Notifications
You must be signed in to change notification settings - Fork 323
/
KafkaDataProvider.cs
175 lines (159 loc) · 7.24 KB
/
KafkaDataProvider.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
using Audit.Core;
using Confluent.Kafka;
using System;
using System.Threading;
using System.Threading.Tasks;
namespace Audit.Kafka.Providers
{
/// <summary>
/// Apache Kafka data provider
/// </summary>
public class KafkaDataProvider : KafkaDataProvider<Null>
{
public KafkaDataProvider(ProducerConfig config) : base(config) { }
public KafkaDataProvider(Action<Configuration.IKafkaProviderConfigurator<Null>> config) : base(config) { }
}
/// <summary>
/// Apache Kafka data provider (keyed messages)
/// </summary>
public class KafkaDataProvider<TKey> : AuditDataProvider
{
private readonly ProducerConfig _producerConfig;
private static readonly object _producerLocker = new object();
private readonly ProducerBuilder<TKey, AuditEvent> _producerBuilder;
private IProducer<TKey, AuditEvent> _producer;
/// <summary>
/// Kafka Topic selector to be used. Default is "audit-topic".
/// </summary>
public Setting<string> Topic { get; set; }
/// <summary>
/// Partition selector to be used as a function of the audit event. Default or NULL means any partition
/// </summary>
public Setting<int?> Partition { get; set; }
/// <summary>
/// Key selector. Optional to use keyed messages. Return the key to be used for a given audit event.
/// </summary>
public Func<AuditEvent, TKey> KeySelector { get; set; }
/// <summary>
/// Key serializer. Optional when using keyed messages and a custom serializer for the key is needed.
/// </summary>
public ISerializer<TKey> KeySerializer { get; set; }
/// <summary>
/// Custom AuditEvent serializer. By default, the audit event is JSON serialized + UTF8 encoded.
/// </summary>
public ISerializer<AuditEvent> AuditEventSerializer { get; set; }
/// <summary>
/// Gets or sets the result handler action. An action to be called for each kafka response
/// </summary>
/// <value>An action to be called for each kafka response.</value>
public Action<DeliveryResult<TKey, AuditEvent>> ResultHandler { get; set; }
/// <summary>
/// Gets or sets the producer builder action. An action to be called before building the producer.
/// </summary>
/// <value>The producer builder extra configuration.</value>
public Action<ProducerBuilder<TKey, AuditEvent>> ProducerBuilderAction { get; set; }
/// <summary>
/// Initializes a new instance of the <see cref="KafkaDataProvider{TKey}"/> class.
/// </summary>
/// <param name="config">The configuration fluent API.</param>
public KafkaDataProvider(Action<Configuration.IKafkaProviderConfigurator<TKey>> config)
{
var kafkaConfig = new Configuration.KafkaProviderConfigurator<TKey>();
if (config != null)
{
config.Invoke(kafkaConfig);
_producerConfig = kafkaConfig._producerConfig;
_producerBuilder = new ProducerBuilder<TKey, AuditEvent>(_producerConfig);
Topic = kafkaConfig._topic;
Partition = kafkaConfig._partition;
KeySelector = kafkaConfig._keySelector;
KeySerializer = kafkaConfig._keySerializer;
AuditEventSerializer = kafkaConfig._auditEventSerializer;
ResultHandler = kafkaConfig._resultHandler;
ProducerBuilderAction = kafkaConfig._producerBuilderAction;
}
}
/// <summary>
/// Initializes a new instance of the <see cref="KafkaDataProvider{TKey}"/> class.
/// </summary>
/// <param name="producerConfig">The producer configuration.</param>
public KafkaDataProvider(ProducerConfig producerConfig)
{
_producerConfig = producerConfig;
_producerBuilder = new ProducerBuilder<TKey, AuditEvent>(_producerConfig);
}
private void EnsureProducer()
{
if (_producer == null)
{
lock(_producerLocker)
{
if (_producer == null)
{
if (KeySerializer != null)
{
_producerBuilder.SetKeySerializer(KeySerializer);
}
_producerBuilder.SetValueSerializer(AuditEventSerializer ?? new DefaultJsonSerializer<AuditEvent>());
// allow extra configuration from the client
ProducerBuilderAction?.Invoke(_producerBuilder);
_producer = _producerBuilder.Build();
}
}
}
}
public override object InsertEvent(AuditEvent auditEvent)
{
EnsureProducer();
return Produce(auditEvent);
}
public override async Task<object> InsertEventAsync(AuditEvent auditEvent, CancellationToken cancellationToken = default)
{
EnsureProducer();
return await ProduceAsync(auditEvent, cancellationToken);
}
public override void ReplaceEvent(object eventId, AuditEvent auditEvent)
{
EnsureProducer();
Produce(auditEvent);
}
public override async Task ReplaceEventAsync(object eventId, AuditEvent auditEvent, CancellationToken cancellationToken = default)
{
EnsureProducer();
await ProduceAsync(auditEvent, cancellationToken);
}
public override T GetEvent<T>(object eventId)
{
throw new NotImplementedException();
}
public override Task<T> GetEventAsync<T>(object eventId, CancellationToken cancellationToken = default)
{
throw new NotImplementedException();
}
private TopicPartition GetTopicPartition(AuditEvent auditEvent)
{
var topic = Topic.GetValue(auditEvent) ?? "audit-topic";
var partitionIndex = Partition.GetValue(auditEvent);
var partition = partitionIndex.HasValue ? new Partition(partitionIndex.Value) : Confluent.Kafka.Partition.Any;
return new TopicPartition(topic, partition);
}
private TKey Produce(AuditEvent auditEvent)
{
var key = KeySelector == null ? default : KeySelector.Invoke(auditEvent);
var message = new Message<TKey, AuditEvent> { Key = key, Value = auditEvent };
var topic = GetTopicPartition(auditEvent);
var result = _producer.ProduceAsync(topic, message).GetAwaiter().GetResult();
ResultHandler?.Invoke(result);
return result.Key;
}
private async Task<TKey> ProduceAsync(AuditEvent auditEvent, CancellationToken cancellationToken)
{
var key = KeySelector == null ? default : KeySelector.Invoke(auditEvent);
var message = new Message<TKey, AuditEvent> { Key = key, Value = auditEvent };
var topic = GetTopicPartition(auditEvent);
var result = await _producer.ProduceAsync(topic, message, cancellationToken);
ResultHandler?.Invoke(result);
return result.Key;
}
}
}