Post Syndicated from Kalyan Janaki original https://aws.amazon.com/blogs/big-data/build-an-end-to-end-change-data-capture-with-amazon-msk-connect-and-aws-glue-schema-registry/
The value of data is time sensitive. Real-time processing makes data-driven decisions accurate and actionable in seconds or minutes instead of hours or days. Change data capture (CDC) refers to the process of identifying and capturing changes made to data in a database and then delivering those changes in real time to a downstream system. Capturing every change from transactions in a source database and moving them to the target in real time keeps the systems synchronized, and helps with real-time analytics use cases and zero-downtime database migrations. The following are a few benefits of CDC:
- It eliminates the need for bulk load updating and inconvenient batch windows by enabling incremental loading or real-time streaming of data changes into your target repository.
- It ensures that data in multiple systems stays in sync. This is especially important if you’re making time-sensitive decisions in a high-velocity data environment.
Kafka Connect is an open-source component of Apache Kafka that works as a centralized data hub for simple data integration between databases, key-value stores, search indexes, and file systems. The AWS Glue Schema Registry allows you to centrally discover, control, and evolve data stream schemas. Kafka Connect and Schema Registry integrate to capture schema information from connectors. Kafka Connect provides a mechanism for converting data from the internal data types used by Kafka Connect to data types represented as Avro, Protobuf, or JSON Schema. AvroConverter, ProtobufConverter, and JsonSchemaConverter automatically register schemas generated by Kafka connectors (source) that produce data to Kafka. Connectors (sink) that consume data from Kafka receive schema information in addition to the data for each message. This allows sink connectors to know the structure of the data to provide capabilities like maintaining a database table schema in a data catalog.
The post demonstrates how to build an end-to-end CDC using Amazon MSK Connect, an AWS managed service to deploy and run Kafka Connect applications and AWS Glue Schema Registry, which allows you to centrally discover, control, and evolve data stream schemas.
Solution overview

On the producer side, for this example we choose a MySQL-compatible Amazon Aurora database as the data source, and we have a Debezium MySQL connector to perform CDC. The Debezium connector continuously monitors the databases and pushes row-level changes to a Kafka topic. The connector fetches the schema from the database to serialize the records into a binary form. If the schema doesn’t already exist in the registry, the schema will be registered. If the schema exists but the serializer is using a new version, the schema registry checks the compatibility mode of the schema before updating the schema. In this solution, we use backward compatibility mode. The schema registry returns an error if a new version of the schema is not backward compatible, and we can configure Kafka Connect to send incompatible messages to the dead-letter queue.
On the consumer side, we use an Amazon Simple Storage Service (Amazon S3) sink connector to deserialize the record and store changes to Amazon S3. We build and deploy the Debezium connector and the Amazon S3 sink using MSK Connect.
Example schema
For this post, we use the following schema as the first version of the table:
Prerequisites
Before configuring the MSK producer and consumer connectors, we need to first set up a data source, MSK cluster, and new schema registry. We provide an AWS CloudFormation template to generate the supporting resources needed for the solution:
- A MySQL-compatible Aurora database as the data source. To perform CDC, we turn on binary logging in the DB cluster parameter group.
- An MSK cluster. To simplify the network connection, we use the same VPC for the Aurora database and the MSK cluster.
- Two schema registries to handle schemas for message key and message value.
- One S3 bucket as the data sink.
- MSK Connect plugins and worker configuration needed for this demo.
- One Amazon Elastic Compute Cloud (Amazon EC2) instance to run database commands.
To set up resources in your AWS account, complete the following steps in an AWS Region that supports Amazon MSK, MSK Connect, and the AWS Glue Schema Registry:
- Choose Launch Stack:

- Choose Next.
- For Stack name, enter suitable name.
- For Database Password, enter the password you want for the database user.
- Keep other values as default.
- Choose Next.
- On the next page, choose Next.
- Review the details on the final page and select I acknowledge that AWS CloudFormation might create IAM resources.
- Choose Create stack.
Custom plugin for the source and destination connector
A custom plugin is a set of JAR files that contain the implementation of one or more connectors, transforms, or converters. Amazon MSK will install the plugin on the workers of the MSK Connect cluster where the connector is running. As part of this demo, for the source connector we use open-source Debezium MySQL connector JARs, and for the destination connector we use the Confluent community licensed Amazon S3 sink connector JARs. Both the plugins are also added with libraries for Avro Serializers and Deserializers of the AWS Glue Schema Registry. These custom plugins are already created as part of the CloudFormation template deployed in the previous step.
Use the AWS Glue Schema Registry with the Debezium connector on MSK Connect as the MSK producer
We first deploy the source connector using the Debezium MySQL plugin to stream data from an Amazon Aurora MySQL-Compatible Edition database to Amazon MSK. Complete the following steps:
- On the Amazon MSK console, in the navigation pane, under MSK Connect, choose Connectors.
- Choose Create connector.
- Choose Use existing custom plugin and then pick the custom plugin with name starting
msk-blog-debezium-source-plugin. - Choose Next.
- Enter a suitable name like
debezium-mysql-connectorand an optional description. - For Apache Kafka cluster, choose MSK cluster and choose the cluster created by the CloudFormation template.
- In Connector configuration, delete the default values and use the following configuration key-value pairs and with the appropriate values:
- name – The name used for the connector.
- database.hostsname – The CloudFormation output for Database Endpoint.
- database.user and database.password – The parameters passed in the CloudFormation template.
- database.history.kafka.bootstrap.servers – The CloudFormation output for Kafka Bootstrap.
- key.converter.region and value.converter.region – Your Region.
Some of these settings are generic and should be specified for any connector. For example:
- connector.class is the Java class of the connector
- tasks.max is the maximum number of tasks that should be created for this connector
Some settings (database.*, transforms.*) are specific to the Debezium MySQL connector. Refer to Debezium MySQL Source Connector Configuration Properties for more information.
Some settings (key.converter.* and value.converter.*) are specific to the Schema Registry. We use the AWSKafkaAvroConverter from the AWS Glue Schema Registry Library as the format converter. To configure AWSKafkaAvroConverter, we use the value of the string constant properties in the AWSSchemaRegistryConstants class:
key.converterandvalue.convertercontrol the format of the data that will be written to Kafka for source connectors or read from Kafka for sink connectors. We useAWSKafkaAvroConverterfor Avro format.key.converter.registry.nameandvalue.converter.registry.namedefine which schema registry to use.key.converter.compatibilityandvalue.converter.compatibilitydefine the compatibility model.
Refer to Using Kafka Connect with AWS Glue Schema Registry for more information.
- Next, we configure Connector capacity. We can choose Provisioned and leave other properties as default
- For Worker configuration, choose the custom worker configuration with name starting
msk-gsr-blogcreated as part of the CloudFormation template. - For Access permissions, use the AWS Identity and Access Management (IAM) role generated by the CloudFormation template
MSKConnectRole. - Choose Next.
- For Security, choose the defaults.
- Choose Next.
- For Log delivery, select Deliver to Amazon CloudWatch Logs and browse for the log group created by the CloudFormation template (
msk-connector-logs). - Choose Next.
- Review the settings and choose Create connector.
After a few minutes, the connector changes to running status.
Use the AWS Glue Schema Registry with the Confluent S3 sink connector running on MSK Connect as the MSK consumer
We deploy the sink connector using the Confluent S3 sink plugin to stream data from Amazon MSK to Amazon S3. Complete the following steps:
-
- On the Amazon MSK console, in the navigation pane, under MSK Connect, choose Connectors.
- Choose Create connector.
- Choose Use existing custom plugin and choose the custom plugin with name starting
msk-blog-S3sink-plugin. - Choose Next.
- Enter a suitable name like
s3-sink-connectorand an optional description. - For Apache Kafka cluster, choose MSK cluster and select the cluster created by the CloudFormation template.
- In Connector configuration, delete the default values provided and use the following configuration key-value pairs with appropriate values:
-
- name – The same name used for the connector.
- s3.bucket.name – The CloudFormation output for Bucket Name.
- s3.region, key.converter.region, and value.converter.region – Your Region.
-
- Next, we configure Connector capacity. We can choose Provisioned and leave other properties as default
- For Worker configuration, choose the custom worker configuration with name starting
msk-gsr-blogcreated as part of the CloudFormation template. - For Access permissions, use the IAM role generated by the CloudFormation template
MSKConnectRole. - Choose Next.
- For Security, choose the defaults.
- Choose Next.
- For Log delivery, select Deliver to Amazon CloudWatch Logs and browse for the log group created by the CloudFormation template
msk-connector-logs. - Choose Next.
- Review the settings and choose Create connector.
After a few minutes, the connector is running.
Test the end-to-end CDC log stream
Now that both the Debezium and S3 sink connectors are up and running, complete the following steps to test the end-to-end CDC:
- On the Amazon EC2 console, navigate to the Security groups page.
- Select the security group
ClientInstanceSecurityGroupand choose Edit inbound rules. - Add an inbound rule allowing SSH connection from your local network.
- On the Instances page, select the instance
ClientInstanceand choose Connect. - On the EC2 Instance Connect tab, choose Connect.
- Ensure your current working directory is
/home/ec2-userand it has the filescreate_table.sql,alter_table.sql,initial_insert.sql, andinsert_data_with_new_column.sql. - Create a table in your MySQL database by running the following command (provide the database host name from the CloudFormation template outputs):
- When prompted for a password, enter the password from the CloudFormation template parameters.
- Insert some sample data into the table with the following command:
- When prompted for a password, enter the password from the CloudFormation template parameters.
- On the AWS Glue console, choose Schema registries in the navigation pane, then choose Schemas.
- Navigate to
db1.sampledatabase.moviesversion 1 to check the new schema created for the movies table:
A separate S3 folder is created for each partition of the Kafka topic, and data for the topic is written in that folder.
- On the Amazon S3 console, check for data written in Parquet format in the folder for your Kafka topic.
Schema evolution
After the initial schema is defined, applications may need to evolve it over time. When this happens, it’s critical for the downstream consumers to be able to handle data encoded with both the old and the new schema seamlessly. Compatibility modes allow you to control how schemas can or can’t evolve over time. These modes form the contract between applications producing and consuming data. For detailed information about different compatibility modes available in the AWS Glue Schema Registry, refer to AWS Glue Schema Registry. In our example, we use backward combability to ensure consumers can read both the current and previous schema versions. Complete the following steps:
- Add a new column to the table by running the following command:
- Insert new data into the table by running the following command:
- On the AWS Glue console, choose Schema registries in the navigation pane, then choose Schemas.
- Navigate to the schema
db1.sampledatabase.moviesversion 2 to check the new version of the schema created for the movies table movies including the country column that you added:
- On the Amazon S3 console, check for data written in Parquet format in the folder for the Kafka topic.
Clean up
To help prevent unwanted charges to your AWS account, delete the AWS resources that you used in this post:
- On the Amazon S3 console, navigate to the S3 bucket created by the CloudFormation template.
- Select all files and folders and choose Delete.
- Enter permanently delete as directed and choose Delete objects.
- On the AWS CloudFormation console, delete the stack you created.
- Wait for the stack status to change to DELETE_COMPLETE.
Conclusion
This post demonstrated how to use Amazon MSK, MSK Connect, and the AWS Glue Schema Registry to build a CDC log stream and evolve schemas for data streams as business needs change. You can apply this architecture pattern to other data sources with different Kafka connecters. For more information, refer to the MSK Connect examples.
About the Author
Kalyan Janaki is Senior Big Data & Analytics Specialist with Amazon Web Services. He helps customers architect and build highly scalable, performant, and secure cloud-based solutions on AWS.





Flora Wu is a Sr. Resident Architect at AWS Data Lab. She helps enterprise customers create data analytics strategies and build solutions to accelerate their businesses outcomes. In her spare time, she enjoys playing tennis, dancing salsa, and traveling.
Daniel Li is a Sr. Solutions Architect at Amazon Web Services. He focuses on helping customers develop, adopt, and implement cloud services and strategy. When not working, he likes spending time outdoors with his family.
Kachi Odoemene is an Applied Scientist at AWS AI. He builds AI/ML solutions to solve business problems for AWS customers.
Taylor McNally is a Deep Learning Architect at Amazon Machine Learning Solutions Lab. He helps customers from various industries build solutions leveraging AI/ML on AWS. He enjoys a good cup of coffee, the outdoors, and time with his family and energetic dog.
Austin Welch is a Data Scientist in the Amazon ML Solutions Lab. He develops custom deep learning models to help AWS public sector customers accelerate their AI and cloud adoption. In his spare time, he enjoys reading, traveling, and jiu-jitsu.


AWS and Hugging Face collaborate to make generative AI more accessible and cost-efficient – This previous week, we announced an expanded collaboration between AWS and 
AWS Pi Day – Join me on March 14 for the third annual 





























Dhiraj Thakur is a Solutions Architect with Amazon Web Services. He works with AWS customers and partners to provide guidance on enterprise cloud adoption, migration, and strategy. He is passionate about technology and enjoys building and experimenting in the analytics and AI/ML space.
Rajdip Chaudhuri is Solutions Architect with Amazon Web Services specializing in data and analytics. He enjoys working with AWS customers and partners on data and analytics requirements. In his spare time, he enjoys soccer.












Sandeep Adwankar is a Senior Technical Product Manager at AWS. Based in the California Bay Area, he works with customers around the globe to translate business and technical requirements into products that enable customers to improve how they manage, secure, and access data.
Srividya Parthasarathy is a Senior Big Data Architect on the AWS Lake Formation team. She enjoys building data mesh solutions and sharing them with the community.

Olivia Michele is a Data Scientist Lead at Ruparupa, where she has worked in a variety of data roles over the past 5 years, including building and integrating Ruparupa data systems with AWS to improve user experience with data and reporting tools. She is passionate about turning raw information into valuable actionable insights and delivering value to the company.
Dariswan Janweri P. is a Data Engineer at Ruparupa. He considers challenges or problems as interesting riddles and finds satisfaction in solving them, and even more satisfaction by being able to help his colleagues and friends, “two birds one stone.” He is excited to be a major player in Indonesia’s technology transformation.
Adrianus Budiardjo Kurnadi is a Senior Solutions Architect at Amazon Web Services Indonesia. He has a strong passion for databases and machine learning, and works closely with the Indonesian machine learning community to introduce them to various AWS Machine Learning services. In his spare time, he enjoys singing in a choir, reading, and playing with his two children.
Nico Anandito is an Analytics Specialist Solutions Architect at Amazon Web Services Indonesia. He has years of experience working in data integration, data warehouses, and big data implementation in multiple industries. He is certified in AWS data analytics and holds a master’s degree in the data management field of computer science.













SaiKiran Reddy Aenugu is a Data Architect in the Amazon Web Services (AWS) Data Lab. He has 10 years of experience implementing data loading, transformation, and visualization processes. SaiKiran currently helps organizations in North America to adopt modern data architectures such as data lakes and data mesh. He has experience in the retail, airline, and finance sectors.
Narendra Merla is a Data Architect in the Amazon Web Services (AWS) Data Lab. He has 12 years of experience in designing and productionalizing both real-time and batch-oriented data pipelines and building data lakes on both cloud and on-premises environments. Narendra currently helps organizations in North America to build and design robust data architectures, and has experience in the telecom and finance sectors.


Unzip the 












Siva Manickam is the Director of Enterprise Architecture, Integrations, Digital Research & Development at Vyaire Medical Inc. In this role, Mr. Manickam is responsible for the company’s corporate functions (Enterprise Architecture, Enterprise Integrations, Data Engineering) and produce function (Digital Innovation Research and Development).
Prahalathan M is the Data Integration Architect at Vyaire Medical Inc. In this role, he is responsible for end-to-end enterprise solutions design, architecture, and modernization of integrations and data platforms using AWS cloud-native services.
Deenbandhu Prasad is a Senior Analytics Specialist at AWS, specializing in big data services. He is passionate about helping customers build modern data architecture on the AWS Cloud. He has helped customers of all sizes implement data management, data warehouse, and data lake solutions.




Subhro Bose is a Senior Data Architect in Emergent Technologies and Intelligence Platform in Amazon. He loves solving science problems with emergent technologies such as AI/ML, big data, quantum, and more to help businesses across different industry verticals succeed within their innovation journey. In his spare time, he enjoys playing table tennis, learn theories of environmental economics and explore the best muffins across the city.
Ketan Karalkar is a Big Data Solutions Consultant at AWS. He has nearly 2 decades of experience helping customers design and build data analytics, and database solutions. He believes in using technology as an enabler to solve real life business problems.
Eva Fang is a Data Scientist within Professional Services in AWS. She is passionate about using the technology to provide value to customers and achieve business outcomes. She is based in London, in her spare time, she likes to watch movies and musicals.
Payal Singh is a Partner Solutions Architect at Amazon Web Services, focused on the Serverless platform. She is responsible for helping partner and customers modernize and migrate their applications to AWS.














Igor Alekseev is a Senior Partner Solution Architect at AWS in Data and Analytics domain. In his role Igor is working with strategic partners helping them build complex, AWS-optimized architectures. Prior joining AWS, as a Data/Solution Architect he implemented many projects in Big Data domain, including several data lakes in Hadoop ecosystem. As a Data Engineer he was involved in applying AI/ML to fraud detection and office automation.
Avinash Kolluri is a Senior Solutions Architect at AWS. He works across Amazon Alexa and Devices to architect and design modern distributed solutions. His passion is to build cost-effective and highly scalable solutions on AWS. In his spare time, he enjoys cooking fusion recipes and traveling.
Vipul Verma is a Sr.Software Engineer at Amazon.com. He has been with Amazon since 2015,solving real-world challenges through technology that directly impact and improve the life of Amazon customers. In his spare time, he enjoys hiking.









Praveen Allam is a Solutions Architect at AWS. He helps customers design scalable, better cost-perfromant enterprise-grade applications using the AWS Cloud. He builds solutions to help organizations make data-driven decisions.
Vivek Singh is Senior Solutions Architect with the AWS Data Lab team. He helps customers unblock their data journey on the AWS ecosystem. His interest areas are data pipeline automation, data quality and data governance, data lakes, and lake house architectures.



































Akira Ajisaka is a Senior Software Development Engineer on the AWS Glue team. He likes open-source software and distributed systems. In his spare time, he enjoys playing both arcade and console games.
Noritaka Sekiyama is a Principal Big Data Architect on the AWS Glue team. He is responsible for building software artifacts to help customers. In his spare time, he enjoys cycling with his new road bike.
Savio Dsouza is a Software Development Manager on the AWS Glue team. His teams work on building and innovating in distributed compute systems and frameworks, namely on Apache Spark.




























