Welcome back! In this second part of my two-part series covering the new Change Event Streaming (CES) feature in SQL Server 2025, I’ll show you how to consume events generated by CES. In Part 1, we provisioned an Azure event hub, generated a SAS token to access the event hub, created the CesDemo sample database, and enabled Change Event Streaming (CES) on the database. We then added tables to an event stream group with deliberate choices for @include_old_values and @include_all_columns. So at this point, CES is now emitting DML changes (inserts, updates, and deletes) from those tables into the event hub.
Note: This post is based on SQL Server 2025 CTP 2.1. Syntax and behavior are subject to subtle changes by the time the product is released. Change Event Streaming (CES) will ultimately be supported across all SKUs of SQL Server, including SQL Server 2025 for Windows, SQL Server 2025 for Linux, Azure SQL Database, and Managed Instance.
Now we’re ready to build a client application to consume generated events. But before we start coding, let’s establish some context so the steps make sense.
First, CES merely writes into Event Hubs. It doesn’t know (or care) who’s listening. It’s up to your client application(s) to subsequently consume those events. Our sample C# application will use the Event Hubs client SDK (specifically, EventProcessorClient) to listen for events.
Every CES client needs somewhere to record progress, as it processes events. This is called a checkpoint, which works like a “bookmark”. Using checkpoints, client applications can stop and later resume where they left off, and not reprocess events that have already been processed. The SDK uses Azure Blob Storage for this purpose.
You’ll also encounter the term consumer group. Think of a consumer group as a “view” of the stream with its own checkpoint. By utilizing multiple consumer groups (one per client application), each application can maintain its own checkpoint for bookmarking its place in the event stream. The Basic tier allows for only one consumer group. Moving to (and paying for) a higher tier than Basic will allow you to manage multiple client applications that consume events simultaneously from the same event hub, each at their own pace, without stepping on each other.
Create a Blob Storage Container
You’ll need a blob container in Azure Storage so that the Event Hubs client SDK can manage checkpoints for your consumer groups.
Create a Storage Account
A blob container lives within a storage account. To create a new storage account:
- In the Azure portal, create a new resource.
- From the Marketplace, create a new Storage Account resource.
- Provide a name for a new storage account in either a new or existing resource group (dashes not permitted).
- For the Primary service, choose Azure Blob Storage or Azure Data Lake Storage Gen 2.
- For Redundancy, choose Locally-redundant storage (LRS) (sufficient for development and testing).
- Click Review + create, and then Create.

Create a Blob Container
Now you can create a new blob container within the new storage account:
- Under Data Storage on the left, click Containers.
- Click + Add container.
- Provide a name for the new container.
- Click Create.

Now get the connection string for the storage account:
- Under Security + Networking on the left, click Access Keys.
- Click Show under the Connection String for key1.
- Click the Copy icon to copy the connection string to the clipboard.
- Paste the connection string into Notepad; it will be needed for the client application configuration.

Create the Visual Studio Project
Alright, we’re ready to roll. We’ll build our consumer client as a simple console app, keeping the the focus on wiring up the stream, deserializing events, and showing what’s happening.
Note: CES consumers can also be built with Azure Functions (I’ll cover that in a later post). Azure Functions hide much of the boilerplate with an Event Hubs trigger, run serverlessly, and scale out automatically. In contrast, building a client “manually” as we’re doing here, gives you maximum control over connection behavior, batching, retry policies, and diagnostics.
Let’s get started!
Launch Visual Studio 2022. Then select Create a new project and choose Console App (C#). Name the project CESClient, click Next, and then click Create.
Install NuGet Packages
First, we’ll need three NuGet packages to support our application. Right-click the CESClient project and choose Manage NuGet Packages. Click the Browse tab, and then locate and install the following packages:
Azure.Messaging.EventHubs.Processor- Includes the Event Hubs client and processor, as well as Azure Blob Storage for checkpoint support.
Microsoft.Extensions.Configuration.Json- Supports external configuration in
appsettings.jsonrather than using hard-coded configuration.
- Supports external configuration in
Newtonsoft.Json- Allows us to deserialize the CloudEvent payload received from the event hub, which is supplied as JSON.
Add a Configuration File
Now create the appsettings.json file where we’ll keep our configuration. This includes connection details and secrets (SAS token, Blob connection string).
- Right-click the project and choose Add > New Item
- Name the file appsettings.json.
- Replace its content with:
{
"EventHub": {
"HostName": "ces-namespace.servicebus.windows.net",
"Name": "ces-hub",
"SasToken": "paste-your-sas-token-here"
},
"BlobStorage": {
"ConnectionString": "paste-your-blob-connection-string-here",
"ContainerName": "ces-checkpoint"
}
}
- For the EventHub property, note the HostName property specifies our event hub namespace name
ces-namespaceas the host name prefix, the Name property specifies our event hub nameces-hub, and the SasToken property holds the SAS token generated for accessing the event hub. All three of these values were established during setup and configuration in Part 1. - For the BlobStorage property, paste in values for the ConnectionString and ContainerName for the Azure Storage blob container that you just created.
- To ensure this file gets copied to the output directory when we build the project, click
appsettings.jsonin the Solution Explorer panel. Then, in the Properties panel set Copy to Output Directory to Copy if newer.
Add the Code
Now supply the following code in Program.cs:
using Azure;
using Azure.Messaging.EventHubs;
using Azure.Messaging.EventHubs.Processor;
using Azure.Messaging.EventHubs.Consumer;
using Azure.Storage.Blobs;
using Microsoft.Extensions.Configuration;
using System;
using System.Collections.Generic;
using System.IO;
using System.Text.Json;
using System.Threading.Tasks;
namespace CESClient
{
public class Program
{
private static int _eventCount;
// Add methods here
}
}
This imports all the namespaces we’ll be referencing and defines a private field as a simple event counter that we’ll increment with each received event.
Next plug in the Main method:
public static async Task Main(string[] args)
{
// Say hello
Console.WriteLine("SQL Server 2025 Change Event Streaming Client");
Console.WriteLine();
Console.Write("Initializing... ");
// Load configuration from appsettings.json
var config = new ConfigurationBuilder()
.SetBasePath(Directory.GetCurrentDirectory())
.AddJsonFile("appsettings.json", optional: false, reloadOnChange: true)
.Build();
// Create a blob container client that the event processor will use for checkpointing
var blobStorageConnectionString = config["BlobStorage:ConnectionString"];
var blobStorageContainerName = config["BlobStorage:ContainerName"];
var storageClient = new BlobContainerClient(blobStorageConnectionString, blobStorageContainerName);
// Create an event processor client to process events in the event hub
var eventHubHostName = config["EventHub:HostName"];
var eventHubName = config["EventHub:Name"];
var sasToken = config["EventHub:SasToken"];
var processor = new EventProcessorClient(
storageClient, // checkpoint store
EventHubConsumerClient.DefaultConsumerGroupName, // Basic tier: one consumer group (e.g., $Default)
eventHubHostName,
eventHubName,
new AzureSasCredential(sasToken)
);
// Register handlers for processing events and errors
processor.ProcessEventAsync += ProcessEventHandler;
processor.ProcessErrorAsync += ProcessErrorHandler;
// Start listening for events
Console.Write("starting... ");
_eventCount = 0;
await processor.StartProcessingAsync();
Console.WriteLine("waiting... press any key to stop.");
Console.ReadKey(intercept: true);
// Stop listening for events
await processor.StopProcessingAsync();
Console.WriteLine("Stopped");
}
This code loads the configuration from appsettings.json, creates the blob container client (for saving checkpoints), spins up the Event Hubs processor with the default consumer group, and attaches handlers for processing events and errors. Finally, it starts and stops cleanly when the user presses any key.
Process Events
Now add the ProcessEventHandler method. This method first parses the outer CloudEvent envelope, then the inner payload, prints helpful metadata, and routes to the insert/update/delete handlers. Finally (and critically), it updates the checkpoint so restarts will resume from the next event. (I explain the CloudEvent payload structure in Part 1.)
private static async Task ProcessEventHandler(ProcessEventArgs eventArgs)
{
try
{
// Deserialize the event data
using var doc = JsonDocument.Parse(eventArgs.Data.Body.ToArray());
var root = doc.RootElement;
var dataJson = root.GetProperty("data");
using var innerDoc = JsonDocument.Parse(dataJson.GetString());
var data = innerDoc.RootElement;
Console.WriteLine($"Processing event... #{++_eventCount}");
// Deserialize the "current" and "old" fields in the eventrow property of the event data to dictionaries
var operation = root.GetProperty("operation").GetString();
var cols = data.GetProperty("eventsource").GetProperty("cols").EnumerateArray();
var current = JsonSerializer.Deserialize<Dictionary<string, string>>(data.GetProperty("eventrow").GetProperty("current").GetString());
var old = JsonSerializer.Deserialize<Dictionary<string, string>>(data.GetProperty("eventrow").GetProperty("old").GetString());
DisplayEventMetadata(eventArgs, root, data);
switch (operation)
{
case "INS":
ProcessInsert(cols, current);
break;
case "UPD":
ProcessUpdate(cols, current, old);
break;
case "DEL":
ProcessDelete(cols, old);
break;
}
Console.WriteLine();
Console.WriteLine(new string('-', 80));
Console.WriteLine();
// Persist progress so we don't reprocess this event on restart
await eventArgs.UpdateCheckpointAsync();
}
catch (Exception ex)
{
Console.ForegroundColor = ConsoleColor.Red;
Console.WriteLine(ex.Message);
Console.ResetColor();
}
}
Display Event Metadata
This method renders a quick “context dump” for each event: it first prints the sequence number and offset from ProcessEventArgs so you can pinpoint the event’s exact position within the event hub (useful for ordering and replay). It then surfaces key CloudEvent fields from the outer envelope; spec/version, the event type, the DML operation (INS, UPD, DEL), timestamp, unique ID, logical ID, and the data content type. Finally, it drills into the inner CES payload to show the database, schema, and table that produced the event. Together, these details make it easy to correlate what you’re seeing in the console with the emitting source and to troubleshoot issues like unexpected operations or schema mismatches.
private static void DisplayEventMetadata(ProcessEventArgs eventArgs, JsonElement root, JsonElement data)
{
Console.WriteLine("Event Args");
Console.WriteLine($" Sequence:Offset => {eventArgs.Data.SequenceNumber}:{eventArgs.Data.Offset}");
Console.WriteLine();
Console.WriteLine("Event Data");
Console.WriteLine($" Spec version: {root.GetProperty("specversion").GetString()}");
Console.WriteLine($" Operation: {root.GetProperty("type").GetString()}");
Console.WriteLine($" Time: {root.GetProperty("time").GetString()}");
Console.WriteLine($" Event ID: {root.GetProperty("id").GetString()}");
Console.WriteLine($" Logical ID: {root.GetProperty("logicalid").GetString()}");
Console.WriteLine($" Operation: {root.GetProperty("operation").GetString()}");
Console.WriteLine($" Data content type: {root.GetProperty("datacontenttype").GetString()}");
Console.WriteLine();
Console.WriteLine("Data");
Console.WriteLine($" Database: {data.GetProperty("eventsource").GetProperty("db").GetString()}");
Console.WriteLine($" Schema: {data.GetProperty("eventsource").GetProperty("schema").GetString()}");
Console.WriteLine($" Table: {data.GetProperty("eventsource").GetProperty("tbl").GetString()}");
Console.WriteLine();
}
Handle Inserts
For inserts, a full “after” image is easiest to read and is a quick way to validate the @include_all_columns setting we established in Part 1.
private static void ProcessInsert(JsonElement.ArrayEnumerator cols, Dictionary<string, string> current)
{
Console.WriteLine("Operation: Insert");
Console.ForegroundColor = ConsoleColor.Green;
foreach (var col in cols)
{
var name = col.GetProperty("name").GetString();
Console.WriteLine($"\t{name}: {current[name]}");
}
Console.ResetColor();
}
Handle Updates
For tables where we’ve enabled @include_old_values in Part 1, you’ll get a great side-by-side view; otherwise you’ll just see the “after” image.
private static void ProcessUpdate(JsonElement.ArrayEnumerator cols, Dictionary<string, string> current, Dictionary<string, string> old)
{
Console.WriteLine("Operation: Update");
foreach (var col in cols)
{
var name = col.GetProperty("name").GetString();
if (old.Count > 0 && current[name] != old[name])
{
Console.ForegroundColor = ConsoleColor.Yellow;
Console.WriteLine($"\t{name}: {current[name]} (old: {old[name]})");
Console.ResetColor();
}
else
{
Console.WriteLine($"\t{name}: {current[name]}");
}
}
}
Handle Deletes
Deletes only have the “before” image, which are useful for auditing and reconciliation purposes.
private static void ProcessDelete(JsonElement.ArrayEnumerator cols, Dictionary<string, string> old)
{
Console.WriteLine("Operation: Delete");
Console.ForegroundColor = ConsoleColor.Red;
foreach (var col in cols)
{
var name = col.GetProperty("name").GetString();
Console.WriteLine($"\t{name}: {old[name]}");
}
Console.ResetColor();
}
Process Errors
Finally, we need to tolerate errors without crashing the application. Should an error occur, this method displays the exception details. Of course, a real-world scenario would require proper error handling; for example, saving the event details to a queue for automatic retry or manual intervention.
private static Task ProcessErrorHandler(ProcessErrorEventArgs e)
{
Console.ForegroundColor = ConsoleColor.Red;
Console.WriteLine(e.Exception.Message);
Console.ResetColor();
return Task.CompletedTask;
}
Run the Application
The moment of truth is here! Go ahead and run the application. The client console window should open, and you should see:
SQL Server 2025 Change Event Streaming Client
Initializing... starting... waiting... press any key to stop.
If errors appear, double-check configuration values and NuGet package installs.
Generate and Monitor Change Events
Let’s exercise inserts, updates, deletes, as well as trigger-driven changes, based on the database schema we setup in Part 1. This will allow us to observe the events being captured in real-time, and examine the CloudEvent payloads as we receive them.
Start SSMS and open a query window to the CesDemo database. Then tile the SSMS and client console windows side-by-side. This way, you can examine the events in the client console window as you generate them from the SSMS window.
Create an Order
In SSMS, run the stored procedure to create a new order:
EXEC CreateOrder @CustomerId = 1
In the client console window, you should observe an INS event for the Order table.
Create Order Details
Now add two details to the order:
EXEC CreateOrderDetail @OrderId = 1, @ProductId = 1, @Quantity = 2
EXEC CreateOrderDetail @OrderId = 1, @ProductId = 2, @Quantity = 1
Expect two corresponding INS events for OrderDetail. And because the OrderDetail trigger adjusts Product.ItemsInStock, also expect UPD events for Product reflecting stock decrements (2 for product 1; 1 for product 2).
Delete an Order
Now call the stored procedure that deletes an entire order, along with the order details.
EXEC DeleteOrder @OrderId = 1
Expect DEL events for the order and its details, and UPD events for Product as stock is restored.
Update a Customer
Now change a customer’s city to Chicago:
UPDATE Customer SET City = 'Chicago' WHERE CustomerId = 1
Expect a UPD event for Customer. Recall that in Part 1, we chose not to include old values for this table, so you’ll see only the “after” values.
Bulk Update Products
Let’s apply a 20% discount on the price of all cameras:
UPDATE Product SET UnitPrice = UnitPrice * 0.8 WHERE Category = 'Camera'
Expect UPD events for each matching row. In Part 1, we included old values and only changed columns for Product, so you’ll see the old and new UnitPrice (and the primary key, which is always included).
Update a Table Without a Primary Key
This last example demonstrates the need for always having a primary key defined on a table:
UPDATE TableWithNoPK SET ItemName = 'Stove' WHERE Id = 3
Expect a UPD event without key columns. So you’ll see that an item name was changed to Stove in some row, but you won’t know which row, rendering this event data as useless information.
Conclusion
That’s a wrap! In this two-part blog post series, you have successfully built out an end-to-end CES pipeline. First you configured SQL Server 2025 to stream changes into Azure Event Hubs with the appropriate SAS credential, event stream group, and table settings. Then you built a C# consumer that reads the CloudEvent-wrapped payloads and uses Azure Blob Storage for checkpoints. Our demo ran with the default consumer group (Basic tier) for a single application, but higher tiers support multiple groups for multiple independent clients. You now have a clean, real-time path from database changes to actionable events!



