ChangeEventConsumer 1.0.2
dotnet add package ChangeEventConsumer --version 1.0.2
NuGet\Install-Package ChangeEventConsumer -Version 1.0.2
<PackageReference Include="ChangeEventConsumer" Version="1.0.2" />
<PackageVersion Include="ChangeEventConsumer" Version="1.0.2" />
<PackageReference Include="ChangeEventConsumer" />
paket add ChangeEventConsumer --version 1.0.2
#r "nuget: ChangeEventConsumer, 1.0.2"
#:package ChangeEventConsumer@1.0.2
#addin nuget:?package=ChangeEventConsumer&version=1.0.2
#tool nuget:?package=ChangeEventConsumer&version=1.0.2
Change Event Consumer Client Library for .NET
Table of Contents
Introduction
The Change Event Consumer is a .NET client library that enables the user to consume Change Events Streaming (CES) events from Azure Event Hubs using a pull-based model. It simplifies working with these Change Events by enabling developers to work with higher-level objects rather than dealing with low-level issues.
This library provides 3 types of consumers that correspond to 3 levels of abstraction:
TransactionChangeEventConsumer - Consumes logical row-level changes grouped in transactions and guarantees ordering within transactions, as well as within an Event Hub partition.
LogicalChangeEventConsumer - Consumes deserialized logical row-level changes, guaranteed ordering within an Event Hub partition.
RawChangeEventConsumer - Consumes raw Cloud Events generated by Change Event Streaming.
Getting Started
Prerequisites
Azure Subscription
You need an active Azure subscription to use Azure services such as Event Hubs.Event Hubs Namespace and Event Hub
You must have an Event Hubs namespace and at least one Event Hub created.Set up Change Event Streaming on Your Database
Ensure that Change Event Streaming is set up so that events are published to the Event Hub.
Install the Package
Install the Change Event Consumer client library via NuGet:
NuGet\Install-Package ChangeEventConsumer
Authenticate the Client
For the Change Event Consumer client library to interact with an Event Hub, it will need to understand how to connect and authorize with it. The easiest means for doing so is to use a connection string.
You will need to get your connection string and Event Hub name from the Event Hub portal.
Key Concepts
Consumer Group
is a view of an entire Event Hub. Consumer groups enable multiple consuming applications to each have a separate view of the event stream, and to read the stream independently at their own pace and from their own position. There can be at most 5 concurrent readers on a partition per consumer group; however it is recommended that there is only one active consumer for a given partition and consumer group pairing. Each active reader receives all of the events from its partition.Partition
is an ordered sequence of events that is held in an Event Hub. Partitions are a means of data organization associated with the parallelism required by event consumers.Partitioning Type
is the type of partitioning used by Change Event Streaming that publishes the data to the Event Hub. The types of partitioning are Default, Column and Table group.Change Event
refers to a single row change published by Change Event Streaming.
API Overview
To consume Change Events:
Iterator:
- Instantiate
TransactionChangeEventConsumer,LogicalChangeEventConsumerorRawChangeEventConsumer. - Call
GetChangeEventIterator(). - Use
ReadNextAsync()to retrieve events.
This way, the user gets a List of ChangeEvent objects that he can iterate through before calling ReadNextAsync again.
Callback:
- Instantiate
TransactionChangeEventConsumer,LogicalChangeEventConsumerorRawChangeEventConsumer. - Call
StartConsumingAsyncwith desired callback.
This way, the user attaches his own callback function that gets called for each event batch.
Change Events are deserialized based on the typesToDeserialize dictionary, mapping fully qualified table names to .NET types.
(This applies only to the TransactionChangeEventConsumer and LogicalChangeEventConsumer)
ChangeEventConsumers
It handles most of the logic of consuming Change Events. You need to create some type of ChangeEventConsumer to consume data from the Event Hub.
Constructor Arguments:
Required
string connectionString
Connection string for the Event Hub.string eventHubName
Name of the Event Hub.Dictionary<string, Type> typesToDeserialize(OnlyTransactionChangeEventConsumerandLogicalChangeEventConsumer)
This dictionary holds key-value pairs that are <Fully qualified name of the database table, .NET type to deserialize this table to>. This changes what type each of your tables are going to be deserialized to. You can specify a default Type for all unspecified Tables to be deserialized to (by using the key "$default"), this can be a Dictionary<string, string> so it acts as a fallback for deserializing any table not listed, avoiding deserialization errors.
For every change event, the consumer will first check to see if the fully qualified table name is specified as a key, and if it is, try to deserialize to the provided Type. If the encountered fully qualified table name is not specified, or deserialization fails, the consumer will try to deserialize to the type with key "$default". If the default key does not exist or deserialization fails, a TableWithoutTypeException is thrown.
Optional
string consumerGroup
Consumer group name the consumer will be in. Defaults to "$Default".PartitioningType partitioningType
Specifies the partitioning type is used for Change Event Streaming. By defaultPartitioningType.DEFAULTis used, which works for all partitioning types. If you are surely not using the Default partitioning for events you are consuming, it is better to specifyPartitioningType.OTHER. There is a boost in consumer performance when other than default partitioning is used andPartitioningType.OTHERis specified.List<string> partitionIds
Specifies the partition id's to consume events from. By default, the consumer will look at all partitions on the Event Hub. When usingDEFAULTpartitioning, specifying partition IDs is not allowed and will result in anArgumentExceptionThis is because when Default partitioning is used in CES, large message chunks and commands of a single transaction can be on different partitions. Specifying partition IDs in that case can lead to the consumer not being able to get full large message/transaction data, breaking the consumer.EventPosition startingPosition
Specifies the position the consumer will start reading from. By default EventPosition.Earliest is used, meaning the consumer will start reading from the earliest event present on the Event Hub.bool raiseTableWithoutTypeException(OnlyTransactionChangeEventConsumerandLogicalChangeEventConsumer)
Specifies whether the consumer will raise an exception or ignore when a Change Event for a Table not intypesToDeserializeis consumed.long deduplicationLimitEvents
Specifies the maximum amount of physical events the deduplicator will keep track of. Setting this value too low can lead to receiving duplicate messages. Setting it too high can result in an out of memory error and crash.List<string> logicalIdSkipList
Specifies logical IDs of events you want to skip. Primarily used in error recovery when a known event causes issues.MessageType messageType
Specifies the CES type of message that we are consuming. Default is CES_DML_V1. CES_DML_V1 is currently the only message type available.
Methods:
AddDeserializationPairAsync(string fullyQualifiedTableName, Type deserializeTo, ChangeDeserializationPairsOptions options = ChangeDeserializationPairsOptions.CONTINUE, EventPosition restartPosition = null)
Dynamically adds the deserialization type for a given table. If the table already has a defined type to deserialize to, returns false and does nothing. For options, you can choose:CONTINUE: Continue consuming from the current position.RESTART: Restart fromrestartPositionorEventPosition.Earliestif the restartPosition in not defined.
Note: when using the Restart option, you will only receive the unconsumed Change Events. This is true unless you are going back a lot of events (>1.000.000), because then you will leave the scope of the deduplication algorithm, and begin receiving already consumed events. This is why defining a deserialization pair for every Table you need is recommended to do before consuming, so you can always use the CONTINUE option, and avoid ambiguity.
UpdateDeserializationPairAsync(string fullyQualifiedTableName, Type deserializeTo, ChangeDeserializationPairsOptions options = ChangeDeserializationPairsOptions.CONTINUE, EventPosition restartPosition = null)
Works the same as AddDeserializationPairAsync, but returns false if the fullyQualifiedTableName does not already exist as a key. Updates the deserialization type and returns true if it exists.StartConsumingAsync(Func<List<T>, Task> onNewEvents, TimeSpan timeSpanPerPartition, CancellationToken cancellationToken, int maximumEventsPerPartition = int.MaxValue)
Starts consuming change events, each partition is polled for the specified time or until the maximum amount of physical events are consumed, a Change Event batch is created, and the specified callback is called for each batch. T is the type of event the user will receive (LogicalTransactionforTransactionChangeEventConsumer,DmlLogRecordforLogicalChangeEventConsumerandSqlCesCloudEventforRawChangeEventConsumer) The cancellation token is used to cancel the consuming process.GetChangeEventIterator()
Returns aChangeEventIteratorassociated with the consumer.
ChangeEventIterator
The ChangeEventIterator is a class that enables the user to consume Change Events with a Change Event Consumer he already has. It does not have a public constructor and is only available to the user through by calling the GetChangeEventIterator() method in the ChangeEventConsumer class.
Methods:
SetTimeoutPeriod(TimeSpan? timeSpan)
Sets the least amount of polling time spent without any new Change Events needed for HasMoreResults() to start returning false.HasMoreResults()
Returns true if timeout period is set to null (default). If timeout period is set and non-null, it returns true if the time spent polling without events is less than the timeout period. Otherwise returns false.ReadNextAsync(TimeSpan timeSpan, int maximumEvents = int.MaxValue)
Polls for events for the given duration or until the max number of Event Hub events is retrieved.
Note: One Change Event may span multiple Event Hub events, or one Event Hub event may trigger many Change Events due to buffering. Returns a list of events (LogicalTransactionforTransactionChangeEventConsumer,DmlLogRecordforLogicalChangeEventConsumerandSqlCesCloudEventforRawChangeEventConsumer)
Events
LogicalTransaction
Represents a transactional change involving multiple rows.
Properties:
string BeginLsnThe begin lsn of the transactionstring CommitLsnThe commit lsn of the transaction.DateTime CommitTimeThe commit time of the transaction.List<DmlLogRecord> DmlLogRecords
The list of single row changes that make the transaction. In the correct order.
DmlLogRecord
Represents a change to a single row. It has all properties and follows the scheme of the raw Change Event Streaming payload.
Some of the more important properties:
string Operation
The type of operation the Change Event is describing ("DEL", "INS", "UPD").Data->EventSource->string Db
The name of the database the table where the change happened is in.Data->EventSource->string SchemaName
The name of the schema the table where the change happened is in.Data->EventSource->string TableName
The name of the table the change happened to.Data->EventRow->Object Old
Describes how the row in the table looked before the change. The type is the one the user specified in the consumer.Data->EventRow->Object Current
Describes how the row in the table looks after the change. The type is the one the user specified in the consumer.
The full hierarchy of properties:
SpecVersion : string
Type : string
Source : string
LogicalId : string
Time : DateTime
DataContentType : string
Operation : string
Data
EventSource
Db : string
Schema : string
Tbl : string
Cols (List)
Name : string
Type : string
Index : string
PkKey (List)
ColumnName : string
Value : string
Transaction
CommitLsn : string
BeginLsn : string
SequenceNumber : int
FinalEvent : bool
CommitTime : DateTime
EventRow
Old : type specified by user-provided dictionary
Current : type specified by user-provided dictionary
SqlCesCloudEvent
Represents a single Event Hub event.
Properties hierarchy:
SpecVersion
Type : string
Source : string
LogicalId : string
Time : string
DataContentType : string
Operation : string
SegmentIndex : int
FinalSegment : bool
Data : byte[]
Example Usage
Using the iterator:
Dictionary<string, Type> typesToDeserialize = new Dictionary<string, Type>();
typesToDeserialize["testdb.testschema.Table"] = typeof(MyClass);
typesToDeserialize["testdb.testschema.Table1"] = typeof(MyClass1);
typesToDeserialize["$default"] = typeof(Dictionary<string, string>);
TransactionChangeEventConsumer newTransactionConsumer = new TransactionChangeEventConsumer("<< CONNECTION STRING >>", " << EVENT HUB NAME >>", typesToDeserialize);
var ei = newTransactionConsumer.GetChangeEventIterator();
while (ei.HasMoreResults())
{
List<LogicalTransaction> events = await ei.ReadNextAsync(TimeSpan.FromSeconds(2));
if(!events.Any())
{
Console.WriteLine("No new changes.");
await Task.Delay(TimeSpan.FromSeconds(5));
}
else
{
foreach (var lt in events)
{
Console.WriteLine("Beginning of transaction.");
foreach (var dml in lt.DmlLogRecords)
PrintDmlLogRecord(dml);
Console.WriteLine("End of transaction.");
}
}
}
void PrintDmlLogRecord(DmlLogRecord dml)
{
Object? eData = dml.Data.EventRow.Current;
if (eData is MyClass myClass)
{
Console.WriteLine("<MyClass> " + dml.Operation + " Event: f1: " + myClass.f1 + ", f2:" + myClass.f2 + ", Time: " + dml.Time);
}
if (eData is MyClass1 myClass1)
{
Console.WriteLine("<MyClass1> " + dml.Operation + " Event: f1: " + myClass1.f1 + ", f2:" + myClass1.f2 + ", Time: " + dml.Time);
}
if (eData is Dictionary<string, string> myDictionary)
{
Console.WriteLine("<Default> Event. Time: " + dml.Time);
}
}
Using the callback:
Transfer event list processing logic to callback function.
async Task onNewEvents(List<LogicalTransaction> events) {
if (!events.Any())
{
Console.WriteLine("No new changes.");
await Task.Delay(TimeSpan.FromSeconds(5));
return;
}
foreach (LogicalTransaction lt in events)
{
Console.WriteLine("TRANSACTION BEGIN");
foreach (var dml in lt.DmlLogRecords)
{
Console.Write("\t");
PrintDmlLogRecord(dml);
}
Console.WriteLine("TRANSACTION END");
}
}
And then specify the callback:
await newTransactionConsumer.StartConsumingAsync(onNewEvents, TimeSpan.FromSeconds(2), cancellationToken: new CancellationToken());
In these examples:
- Events from table "testdb.testschema.Table" are deserialized into
MyClass. - Events from table "testdb.testschema.Table1" are deserialized into
MyClass1. - All other tables use the default type
Dictionary<string, string>. - The consumer polls every 2 seconds.
- It consumes from all partitions using the "$Default" consumer group.
- All events will be grouped into LogicalTransaction events, because the transaction consumer is used.
Creating User Classes
User Classes are classes the user specifies he wants Change Event data to be deserialized to.
When creating the User Classes, special care must be put into correctness of types and names of class properties.
Follow these steps when creating User Classes:
- Make all of your properties public.
- Name your properties the same as column names in the database.
- Use the correct .NET types for your properties according to the data type mapping below.
Data type mapping
Every SQL type currently supported by CES is provided in this mapping. If you have an unsupported data type in your table, CES won't send data about it, so when deserialized it will be null.
If you do not use correct types, a deserialization error might occur!
| SQL types | .NET types |
|---|---|
| uniqueidentifier | Guid |
| tinyint | Byte |
| smallint | Int16 |
| int | Int32 |
| bigint | Int64 |
| bit | Int16 |
| decimal | Decimal |
| numeric | Decimal |
| money | Decimal |
| smallmoney | Decimal |
| float | Double |
| real | Single |
| date | DateTime |
| time | TimeSpan |
| datetime2 | DateTime |
| datetimeoffset | DateTimeOffset |
| datetime | DateTime |
| smalldatetime | DateTime |
| char | String |
| varchar | String |
| nchar | String |
| nvarchar | String |
| binary | String |
| varbinary | String |
Learn more about Target Frameworks and .NET Standard.
-
net8.0
- Azure.Messaging.EventHubs (>= 5.11.6)
- Azure.Messaging.EventHubs.Processor (>= 5.11.6)
- CloudNative.CloudEvents.Avro (>= 2.8.0)
- CloudNative.CloudEvents.NewtonsoftJson (>= 2.8.0)
- Newtonsoft.Json (>= 13.0.3)
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 release.