Kafka Connector Integration

Earn 25 points (50 with Pro) in two steps

  1. ① Read through the lesson — each section gets a ✓ as you scroll through it.
  2. ② When every section has a ✓, tap Complete lesson.

0 of 12 read · keep scrolling

✦ See fewer ads and earn double points — 50 a lesson instead of 25 — with Pro

Module: Maintaining Azure Cosmos DB Solutions

Lesson: Kafka Connector Integration for Data Movement

Introduction: Bridging Streaming Data and NoSQL Storage

In modern distributed systems, data rarely remains static. It flows from sensors, user interactions, and microservices into various storage engines. Apache Kafka has established itself as the industry standard for distributed event streaming, acting as a high-throughput backbone for asynchronous data processing. However, storing this streaming data in a format that allows for low-latency queries and global distribution is a separate challenge. This is where Azure Cosmos DB enters the picture.

The Kafka Connector for Azure Cosmos DB provides a critical bridge between these two worlds. By integrating Kafka with Cosmos DB, you enable a architecture where events flowing through topics are automatically persisted into your NoSQL database, or conversely, changes within your database are streamed back out to Kafka for downstream consumption. Understanding how to configure, maintain, and troubleshoot this connector is essential for any engineer tasked with building data pipelines that require high availability and massive scale.

This lesson explores the mechanics of the Kafka Connector, the configuration patterns required to optimize performance, and the operational best practices to ensure your data movement remains reliable under heavy load.


Not read yet

Understanding the Architecture of the Kafka Connector

The Kafka Connect framework is a tool for streaming data between Apache Kafka and other systems. It consists of two primary types of components: Source Connectors and Sink Connectors. When working with Cosmos DB, you will primarily interact with these two patterns to facilitate bi-directional data flow.

The Sink Connector

The Sink Connector is the most common use case. It consumes data from one or more Kafka topics and writes that data into an Azure Cosmos DB container. This is ideal for scenarios where you need to archive event logs, materialize views of event-sourced data, or power dashboards that require real-time updates. The connector handles the heavy lifting of mapping Kafka messages to JSON documents, managing batch sizes, and retrying failed writes.

The Source Connector

The Source Connector works in the opposite direction. It monitors the Change Feed of an Azure Cosmos DB container and publishes these changes as messages onto a Kafka topic. This is incredibly powerful for event-driven architectures where you want to trigger downstream processes—such as sending an email, updating a search index, or invalidating a cache—whenever a record in your database is created or modified.

Callout: Change Feed vs. Sink Connector It is important to distinguish between the two directions. The Sink Connector is essentially a "Kafka-to-Cosmos" pipeline, while the Source Connector is a "Cosmos-to-Kafka" pipeline. When using the Source Connector, you are essentially exposing your database's internal state changes to the rest of your ecosystem, turning your database into a producer of events rather than just a passive store.


Not read yet

Setting Up the Environment

Before you can move data, you must ensure that your environment is prepared for the connector. The Kafka Connector for Azure Cosmos DB is typically deployed as a plugin within a Kafka Connect cluster.

Prerequisites

  1. Azure Cosmos DB Account: You must have an active SQL API (Core) account.
  2. Kafka Cluster: A running Kafka cluster (this can be on-premises, managed Kafka like Confluent Cloud, or Azure Event Hubs with Kafka support).
  3. Kafka Connect Worker: A set of workers that run the connector tasks.
  4. Connector JAR: You must download the latest version of the Azure Cosmos DB Kafka Connector JAR file from the official repository or Maven Central.

Step-by-Step Installation

  1. Locate the Plugin Directory: On your Kafka Connect worker nodes, identify the plugin.path directory defined in your connect-distributed.properties file.
  2. Copy the JAR: Move the downloaded Cosmos DB connector JAR file into this directory. If there are dependencies, ensure they are present as well.
  3. Restart Workers: You must restart the Kafka Connect workers so that they can discover the new plugin classes.
  4. Verify Discovery: You can verify that the plugin is loaded by sending a GET request to the Kafka Connect REST API: GET /connector-plugins. You should see com.azure.cosmos.kafka.connect.sink.CosmosSinkConnector and the corresponding source connector in the response list.

Not read yet

Configuring the Sink Connector

The Sink Connector requires a JSON configuration file that defines how to map Kafka topics to your Cosmos DB collections. Below is a detailed breakdown of the essential configuration properties.

Essential Configuration Properties

  • connector.class: This identifies the class to run. For the sink, it is com.azure.cosmos.kafka.connect.sink.CosmosSinkConnector.
  • tasks.max: Controls the parallelism. Set this based on the number of partitions in your Kafka topic and the throughput requirements of your Cosmos DB container.
  • topics: A comma-separated list of Kafka topics that the connector should subscribe to.
  • connect.cosmos.connection.endpoint: The URI of your Cosmos DB account.
  • connect.cosmos.connection.key: The primary or secondary key for authentication.
  • connect.cosmos.database.name: The target database name.
  • connect.cosmos.container.name: The target container name.

Example Configuration Snippet

{
  "name": "cosmos-sink-connector",
  "config": {
    "connector.class": "com.azure.cosmos.kafka.connect.sink.CosmosSinkConnector",
    "tasks.max": "3",
    "topics": "user-activity-topic",
    "connect.cosmos.connection.endpoint": "https://your-account.documents.azure.com:443/",
    "connect.cosmos.connection.key": "your-secret-key",
    "connect.cosmos.database.name": "analytics-db",
    "connect.cosmos.container.name": "user-events",
    "connect.cosmos.sink.batch.size": "100",
    "connect.cosmos.sink.bulk.enabled": "true"
  }
}

Note: Always use Azure Key Vault or a similar secret management system to store your connect.cosmos.connection.key. Never hardcode keys in configuration files that are checked into version control.


Not read yet

Optimizing Performance: Bulk Operations and Throughput

One of the most common mistakes when using the Kafka Connector is failing to tune it for the specific throughput requirements of the workload. By default, the connector might not be optimized for high-volume ingestion, leading to bottlenecks in the Kafka Connect cluster.

Enabling Bulk Support

Cosmos DB supports a bulk execution mode that significantly improves throughput by reducing the number of round-trips to the server. When configuring your sink connector, you should always set connect.cosmos.sink.bulk.enabled to true. This allows the connector to group multiple operations into a single request, which is much more efficient for the underlying Cosmos DB SDK.

Batch Size Tuning

The connect.cosmos.sink.batch.size property determines how many messages the connector accumulates before flushing them to Cosmos DB. If you set this too low, you will have high network overhead. If you set it too high, you might hit the request size limit (which is 2MB per request in Cosmos DB). A good starting point is 100 to 500 documents per batch, but you should perform load testing to find the "sweet spot" for your specific document size.

Partitioning Strategy

Your Cosmos DB partition key choice is critical for performance. Ensure that the messages coming from Kafka contain a property that maps to your Cosmos DB partition key. If your partition key is not present in the Kafka message, the connector will fail to write the record, or it will use a default that might lead to "hot partitions" in your database.


Not read yet

Implementing the Source Connector: Streaming from Cosmos DB

The Source Connector is equally important for keeping downstream systems in sync with your database. This is frequently used for microservices that need to react to data changes.

Configuration for Source

The source connector requires the connect.cosmos.source.connector.CosmosSourceConnector class. You must also specify the change feed configuration, such as the starting point for reading changes.

{
  "name": "cosmos-source-connector",
  "config": {
    "connector.class": "com.azure.cosmos.kafka.connect.source.CosmosSourceConnector",
    "tasks.max": "1",
    "connect.cosmos.connection.endpoint": "https://your-account.documents.azure.com:443/",
    "connect.cosmos.connection.key": "your-secret-key",
    "connect.cosmos.database.name": "analytics-db",
    "connect.cosmos.container.name": "user-events",
    "connect.cosmos.source.changefeed.startFromBeginning": "true",
    "kafka.topic": "cosmos-changes-topic"
  }
}

Managing Change Feed Offsets

The Kafka Connect framework tracks the progress of the source connector using offsets. This ensures that if a worker node crashes, it can resume reading the change feed from the exact point where it left off. This mechanism is built-in and highly reliable, provided that your Kafka Connect cluster has a persistent storage location for its offset data.

Callout: Reliability and At-Least-Once Delivery Both the Sink and Source connectors provide at-least-once delivery guarantees. This means that in the event of a network failure or a crash, a message might be processed more than once. Your downstream applications and your Cosmos DB document schemas should be designed to be idempotent—meaning that processing the same message multiple times does not result in incorrect state.


Not read yet

Operational Best Practices

Maintaining a production-grade data pipeline involves more than just getting the initial configuration correct. You must plan for monitoring, scaling, and handling errors.

Monitoring and Observability

  • JMX Metrics: Kafka Connect exposes extensive metrics via JMX. You should monitor request-latency-avg, record-error-rate, and batch-size-avg. These metrics will tell you if the connector is struggling to keep up with the incoming data volume.
  • Azure Monitor: Since the connector interacts directly with Cosmos DB, keep an eye on the Total Request Units (RU) consumption in the Azure portal. If the connector is causing RU spikes, you may need to increase the throughput of your container or optimize the batching settings.
  • Dead Letter Queues (DLQ): Kafka Connect supports DLQs. If a message cannot be processed (e.g., due to a schema mismatch or a serialization error), it can be sent to a dedicated "dead letter" Kafka topic instead of stopping the entire connector. This is a vital feature for maintaining uptime.

Handling Schema Evolution

Data formats change over time. If your Kafka messages are in Avro or JSON format, ensure that you use a Schema Registry. The Kafka Connector can leverage the Schema Registry to validate messages before attempting to write them to Cosmos DB. Without this, a single malformed message could cause the connector to enter a crash loop.

Scaling the Connector

If you find that the lag in your Kafka topics is growing, the first step is to increase the number of Kafka partitions. Since each task in Kafka Connect is assigned a subset of partitions, you can then increase tasks.max to match the partition count. This parallelizes the ingestion process and allows you to consume data at a much higher rate.


Not read yet

Common Pitfalls and How to Avoid Them

Even with a solid understanding of the mechanics, engineers often run into specific, preventable issues.

1. Incompatible Partition Keys

  • The Problem: The most frequent cause of "400 Bad Request" errors is an attempt to insert a document that lacks the required partition key or has a partition key value that does not match the metadata.
  • The Fix: Always validate the schema of your incoming Kafka messages against your Cosmos DB container schema. Use a Kafka Connect Transform (SMT) if you need to modify the message structure before it reaches the sink.

2. Throughput Throttling

  • The Problem: The connector attempts to write data faster than the provisioned throughput of the Cosmos DB container, resulting in 429 Too Many Requests errors.
  • The Fix: While the connector handles retries automatically, excessive throttling indicates a design flaw. Consider using Autoscale throughput on your Cosmos DB container or implementing a more aggressive backoff strategy in your connector configuration.

3. Misconfigured Offset Storage

  • The Problem: The Kafka Connect cluster is configured with an ephemeral offset storage, causing the connector to re-read the entire change feed from the beginning every time a worker restarts.
  • The Fix: Ensure that your connect-distributed.properties file points to a durable Kafka topic for offset.storage.topic. This topic should have a high replication factor and be compacted to prevent data loss.

4. Ignoring the "Max Request Size" Limit

  • The Problem: Attempting to send a single batch that exceeds the 2MB limit allowed by the Cosmos DB SDK.
  • The Fix: Use the connect.cosmos.sink.batch.size property to limit the number of documents per batch, and ensure that your individual document sizes are well under the limit. If you are dealing with large documents, you may need to implement a pre-processing step to split them.

Not read yet

Quick Reference: Connector Configuration Parameters

Parameter Purpose Recommended Action
tasks.max Controls parallelism Set to match Kafka partition count
connect.cosmos.sink.bulk.enabled Enables bulk API Always set to true
connect.cosmos.sink.batch.size Documents per batch Start at 100, tune based on size
errors.tolerance Error handling Set to all or dlq for production
connect.cosmos.source.changefeed.startFromBeginning Change feed offset Set true for historical load

Advanced Integration Patterns

As your system grows, you may find that the standard sink/source patterns are not enough. Here are two advanced scenarios that often arise in production environments.

Pattern 1: Filtering and Transformation with SMTs

Kafka Connect supports Single Message Transforms (SMTs). These are small code snippets that run inside the connector to modify or filter messages before they are processed. For example, if you have a Kafka topic containing multiple types of events, you can use a Filter transform to discard events that aren't relevant to a specific Cosmos DB container. This saves on storage costs and RU consumption.

Pattern 2: Multi-Region Writes

If your Cosmos DB account is globally distributed with multi-region writes enabled, the Kafka Connector can be configured to write to the local region. This reduces latency significantly. To achieve this, ensure your connector is running on infrastructure within the same Azure region as your target Cosmos DB endpoint.


Not read yet

Troubleshooting Checklist

When things go wrong, follow this systematic approach to isolate the issue:

  1. Check Connector Status: Use GET /connectors/{name}/status to see if the connector is in a FAILED state. The status response will often contain the stack trace of the error.
  2. Inspect Kafka Connect Logs: The logs on the worker nodes are your primary source of truth. Look for CosmosException to identify database-level errors.
  3. Verify Authentication: Ensure the connection string and key have not expired and that the Kafka Connect node has network connectivity to the Cosmos DB endpoint (check for firewall rules or VNET restrictions).
  4. Test Connectivity: Use a simple curl command from the worker node to the Cosmos DB endpoint to ensure there are no network-level blocks.
  5. Review Throughput: Check the Cosmos DB metrics in Azure to see if the container is hitting its RU limit during the time of the failure.

Not read yet

Best Practices for Long-Term Maintenance

  1. Version Management: Keep your connector JAR files updated. The Azure team frequently releases updates that include performance improvements and bug fixes for the Cosmos DB SDK.
  2. Automated Deployments: Treat your connector configurations as code. Use a CI/CD pipeline to deploy connector definitions via the Kafka Connect REST API rather than manual configuration.
  3. Graceful Shutdowns: When performing maintenance on your Kafka Connect cluster, ensure you perform a graceful shutdown of the connectors. This allows the connector to finish its current batch and commit its offsets, preventing duplicate processing upon restart.
  4. Capacity Planning: Periodically review the throughput of your Kafka topics. If the volume of data is increasing due to business growth, you must proactively scale the Kafka Connect cluster and the Cosmos DB throughput to avoid data ingestion lag.
  5. Security Audits: Regularly rotate your Cosmos DB keys used by the connector. Use a secret manager that supports programmatic rotation to minimize downtime.

Not read yet

Key Takeaways

  • Connectivity: The Kafka Connector acts as a vital bridge for bi-directional data flow, enabling event-driven architectures that combine the strengths of Kafka's streaming and Cosmos DB's NoSQL storage.
  • Performance Tuning: Bulk operations and appropriate batch sizing are not optional for high-throughput systems; they are mandatory configurations that define your ability to scale.
  • Operational Resilience: Always use Dead Letter Queues and robust offset storage to ensure that your data pipeline is resilient to transient failures and can recover automatically from crashes.
  • Observability: Treat the connector as a first-class application. Monitor JMX metrics and Cosmos DB RU consumption to proactively identify bottlenecks before they impact your end-users.
  • Idempotency: Because the connector guarantees "at-least-once" delivery, design your downstream logic to be idempotent. This is the single most important design principle for building reliable distributed data systems.
  • Scaling: Align your Kafka partitions with your connector tasks.max settings to ensure that data ingestion is evenly distributed across your infrastructure.
  • Configuration Management: Store configurations in version control and use automated deployment processes to maintain consistency across development, staging, and production environments.

By mastering these concepts, you transition from simply "connecting" two services to architecting a reliable, scalable, and maintainable data movement strategy. The Kafka Connector for Azure Cosmos DB is a powerful tool, and when handled with the right operational rigor, it forms the backbone of highly responsive, event-driven applications.

Not read yet

Each section gets a ✓ as you scroll through it. Tap the button to jump to the next one.