Overview
Couchbase Sink Connector consumes records from Kafka topics and writes them to Couchbase. The default handler converts non-null record values to JSON documents and performs full-document upserts; tombstone records delete the corresponding documents. The Connector also supports N1QL, Analytics, and sub-document handlers for conditional updates, analytical data writes, and partial document changes. The Connector can route records by Kafka topic to different Couchbase Buckets, Scopes, and Collections. A document ID can come from a business field in the message value, the Kafka record key, or a combination of topic, partition, and offset. Kafka Connect advances the offset for a batch only after the target write completes, but the Couchbase write and Kafka offset commit do not form a cross-system transaction.Prerequisites
- For the default KV, N1QL, and sub-document handlers, the target Bucket, Scope, and Collection must already exist, and the connection account must have the corresponding read and write permissions; when using the Analytics handler, a usable Couchbase Analytics service and target dataset are also required.
- When using TLS, prepare a trusted PEM CA certificate or Java truststore; when using mutual TLS, also prepare the client certificate bundle and its password.
- The Couchbase node addresses used by
couchbase.seed.nodesmust be accessible from the Connect Worker; when specifying a custom port, use the KV port.
License
Uses Apache License 2.0.Quick Start
Prepare a Connect Cluster, Kafka, the target Couchbase Bucket, and Collection in advance, and confirm network connectivity and access permissions. For preparation and management operations, see Manage Connectors.<topic-name>, <couchbase-host>, <username>, <password>, and <bucket-name> with the actual values. The default handler writes non-null records to the _default._default Collection in the target Bucket; to use another Scope or Collection, add couchbase.default.collection.
Configuration
Kafka Input and Tasks
connector.class
Specifies the Couchbase Sink Connector implementation class to load.
- Type:
string - Default: None
- Importance: High
- Valid values / notes: Use
com.couchbase.connect.kafka.CouchbaseSinkConnector. - Required: Yes
topics
Specifies the list of Kafka topics consumed by the Connector.
- Type:
list - Default: Empty list
- Importance: High
- Valid values / notes: Use comma-separated topic names; exactly one non-empty value must be configured between
topicsandtopics.regex.
topics.regex
Uses a Java regular expression to select the Kafka topics consumed by the Connector.
- Type:
string - Default: Empty string
- Importance: High
- Valid values / notes: Enter a valid Java Pattern; exactly one non-empty value must be configured between
topicsandtopics.regex.
tasks.max
Specifies the maximum number of Tasks that Kafka Connect can start for this Connector.
- Type:
int - Default:
1 - Importance: High
- Valid values / notes: Must be at least
1; effective parallelism is also limited by the number of input partitions and Couchbase processing capacity.
key.converter
Specifies the Connector-level Converter for Kafka record keys.
- Type:
class - Default:
null, inherited from the Worker configuration - Importance: Low
- Valid values / notes: Use an instantiable Kafka Connect Converter; the record key can be used as a fallback source for the Couchbase document ID.
value.converter
Specifies the Connector-level Converter for Kafka record values.
- Type:
class - Default:
null, inherited from the Worker configuration - Importance: Low
- Valid values / notes: Use a Converter compatible with the input data; the conversion result must be serializable as JSON by the Connector.
Couchbase Connection and Authentication
couchbase.seed.nodes
Specifies the seed nodes used to discover the Couchbase cluster.
- Type:
list - Default: None
- Importance: High
- Valid values / notes: Enter one or more accessible node addresses separated by commas; a custom port must be the KV port.
- Required: Yes
couchbase.username
Specifies the Couchbase username used for password authentication.
- Type:
string - Default: None
- Importance: High
- Valid values / notes: Configuration parsing requires this value; when client certificate authentication is used, the value is ignored at runtime.
- Required: Yes
couchbase.password
Specifies the Couchbase password used for password authentication.
- Type:
password - Default: None
- Importance: High
- Valid values / notes: Configuration parsing requires this value; the environment variable
KAFKA_COUCHBASE_PASSWORDcan override it after parsing. When client certificate authentication is used, the value is ignored at runtime. - Required: Yes
couchbase.bucket
Specifies the default target Bucket and the default Bucket for a Collection Keyspace that omits the Bucket.
- Type:
string - Default: Empty string
- Importance: High
- Valid values / notes: Handlers that use KV Collections must configure an existing and accessible Bucket; the Analytics handler may omit it.
couchbase.network
Selects the cluster address network view used by the Couchbase SDK.
- Type:
string - Default:
auto - Importance: Medium
- Valid values / notes: Can be
auto,default, orexternal.
couchbase.bootstrap.timeout
Sets the maximum time to wait for a Couchbase connection during startup.
- Type:
string - Default:
30s - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
ms,s,m,h, ord;0is also allowed.
TLS and Certificates
couchbase.enable.tls
Controls whether to enable TLS for Couchbase connections.
- Type:
boolean - Default:
false - Importance: Medium
- Valid values / notes: Can be
trueorfalse; secure endpoints should enable TLS and configure appropriate trust material.
couchbase.enable.hostname.verification
Controls whether to verify the hostname in the TLS server certificate.
- Type:
boolean - Default:
true - Importance: Medium
- Valid values / notes: Can be
trueorfalse; disabling it weakens server identity verification.
couchbase.trust.store.path
Specifies the absolute path to the Java KeyStore containing trusted Couchbase CAs.
- Type:
string - Default: Empty string
- Importance: Medium
- Valid values / notes: The path must be readable; this is the trust-material alternative to
couchbase.trust.certificate.path.
couchbase.trust.store.password
Specifies the truststore password.
- Type:
password - Default: Empty string
- Importance: Medium
- Valid values / notes: Provide the value required by the truststore; the environment variable
KAFKA_COUCHBASE_TRUST_STORE_PASSWORDcan override it.
couchbase.trust.certificate.path
Specifies the absolute path to a trusted Couchbase CA PEM certificate.
- Type:
string - Default: Empty string
- Importance: Medium
- Valid values / notes: The path must be readable; this is the PEM alternative to
couchbase.trust.store.path.
couchbase.client.certificate.path
Specifies the path to a Java KeyStore or PKCS12 file containing the client private key and certificate chain.
- Type:
string - Default: Empty string
- Importance: Medium
- Valid values / notes: Setting a non-empty value selects client certificate authentication and causes the username and password to be ignored at runtime.
couchbase.client.certificate.password
Specifies the client certificate bundle password.
- Type:
password - Default: Empty string
- Importance: Medium
- Valid values / notes: Provide the value required by the certificate bundle; the environment variable
KAFKA_COUCHBASE_CLIENT_CERTIFICATE_PASSWORDcan override it.
Routing and Document IDs
couchbase.default.collection
Specifies the default target Keyspace and allows per-Kafka-topic overrides.
- Type:
string - Default:
_default._default - Importance: Medium
- Valid values / notes: Use
scope.collectionorbucket.scope.collection; when the Bucket is omitted,couchbase.bucketis used. For a per-topic override, usecouchbase.default.collection[<kafka-topic>]; topics without an override inherit the base value.
couchbase.topic.to.collection
Uses the legacy mapping to specify a target Keyspace by Kafka topic.
- Type:
list - Default: Empty list
- Importance: Medium
- Valid values / notes: Entries use the format
topic=scope.collectionortopic=bucket.scope.collection; prefercouchbase.default.collection[<kafka-topic>]. - Deprecation status: Deprecated
- Replacement:
couchbase.default.collection[<kafka-topic>]
couchbase.document.id
Sets the format used to extract a Couchbase document ID from the record value.
- Type:
string - Default: Empty string
- Importance: Medium
- Valid values / notes: JSON Pointer placeholders such as
${/id}orprefix::${/id}can be used. When the value is empty or extraction fails, the Connector falls back to the Kafka record key and then to the topic, partition, and offset. For a per-topic override, usecouchbase.document.id[<kafka-topic>].
couchbase.topic.to.document.id
Uses the legacy mapping to specify a document ID format by Kafka topic.
- Type:
list - Default: Empty list
- Importance: Medium
- Valid values / notes: Entries use the format
topic=<document-id-format>; each entry can contain only one mapping separator. Prefercouchbase.document.id[<kafka-topic>]. - Deprecation status: Deprecated
- Replacement:
couchbase.document.id[<kafka-topic>]
couchbase.remove.document.id
Controls whether to remove the field used to generate the document ID from the written document.
- Type:
boolean - Default:
false - Importance: Medium
- Valid values / notes: Meaningful only when all JSON Pointers in
couchbase.document.idare extracted successfully. For a per-topic override, usecouchbase.remove.document.id[<kafka-topic>].
couchbase.document.expiration
Sets the document expiration time.
- Type:
string - Default:
0 - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
ms,s,m,h, ord;0means no expiration. For a per-topic override, usecouchbase.document.expiration[<kafka-topic>].
Write Handling and Retries
couchbase.sink.handler
Specifies the SinkHandler implementation that converts Sink records into Couchbase operations.
- Type:
class - Default:
com.couchbase.connect.kafka.handler.sink.UpsertSinkHandler - Importance: Medium
- Valid values / notes: Use a loadable SinkHandler class with an accessible constructor; a
couchbase.document.modeother thanDOCUMENToverrides this value.
couchbase.document.mode
Selects document, sub-document, or N1QL handling through the compatibility mode.
- Type:
string - Default:
DOCUMENT - Importance: Medium
- Valid values / notes: Can be
DOCUMENT,SUBDOCUMENT, orN1QL;N1QLandSUBDOCUMENToverridecouchbase.sink.handler. - Deprecation status: Deprecated
- Replacement:
couchbase.sink.handler
couchbase.retry.timeout
Sets the total time window for the Connector to retry failed writes.
- Type:
string - Default:
0 - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
ms,s,m,h, ord;0means no retry. This setting differs from Kafka Connect’serrors.retry.timeout; on failure, the entire batch may be replayed.
Write Durability
couchbase.durability
Sets the enhanced durability level for Couchbase Server 6.5 or later.
- Type:
string - Default:
NONE - Importance: Medium
- Valid values / notes: Can be
NONE,MAJORITY,MAJORITY_AND_PERSIST_TO_ACTIVE, orPERSIST_TO_MAJORITY. A non-NONEvalue cannot be used together with a non-NONEcouchbase.persist.toorcouchbase.replicate.to.
couchbase.persist.to
Sets the disk-persistence requirement for older Couchbase Server versions.
- Type:
string - Default:
NONE - Importance: Medium
- Valid values / notes: Can be
NONE,ACTIVE,ONE,TWO,THREE, orFOUR, and the cluster topology must be able to satisfy the requirement. For Couchbase Server 6.5 or later, prefercouchbase.durability.
couchbase.replicate.to
Sets the in-memory replica confirmation requirement for older Couchbase Server versions.
- Type:
string - Default:
NONE - Importance: Medium
- Valid values / notes: Can be
NONE,ONE,TWO, orTHREE, and the cluster topology must be able to satisfy the requirement. For Couchbase Server 6.5 or later, prefercouchbase.durability.
N1QL Handler
couchbase.n1ql.operation
Selects the update method used by N1qlSinkHandler.
- Type:
string - Default:
UPDATE - Importance: Medium
- Valid values / notes: Can be
UPDATEorUPDATE_WHERE; applies only toN1qlSinkHandler.
couchbase.n1ql.where.fields
Specifies the fields and optional literals used to construct the UPDATE_WHERE condition.
- Type:
list - Default: Empty list
- Importance: Medium
- Valid values / notes: Use field names or
field:literalentries; a non-empty list is required whenUPDATE_WHEREis selected.
couchbase.n1ql.create.document
Controls whether N1QL UPDATE mode creates a missing document through MERGE.
- Type:
boolean - Default:
true - Importance: Medium
- Valid values / notes: Applies only to
UPDATEmode ofN1qlSinkHandler.
Analytics Handler
couchbase.analytics.max.records.in.batch
Sets the maximum number of records in each Analytics UPSERT or DELETE statement batch.
- Type:
int - Default:
100 - Importance: Medium
- Valid values / notes: ConfigDef accepts any 32-bit integer; in practice, use a positive integer. This setting has not been committed as a stable interface.
couchbase.analytics.max.size.in.batch
Sets the maximum total number of bytes for document data in an Analytics UPSERT batch.
- Type:
string - Default:
5m - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
b,k,m, org; in practice, use a positive size limit. This setting has not been committed as a stable interface.
couchbase.analytics.query.timeout
Sets the client timeout for Analytics query requests.
- Type:
string - Default:
5m - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
ms,s,m,h, ord;0is also allowed. This setting has not been committed as a stable interface.
Sub-Document Handler
couchbase.subdocument.path
Specifies the sub-document path modified by SubDocumentSinkHandler.
- Type:
string - Default: Empty string
- Importance: Medium
- Valid values / notes: Use a non-empty path; an ordinary value represents a fixed path, while a value beginning with
/represents a JSON Pointer that extracts a dynamic path from each message.
couchbase.subdocument.operation
Selects the sub-document mutation operation.
- Type:
string - Default:
UPSERT - Importance: Medium
- Valid values / notes: Can be
UPSERT,ARRAY_PREPEND, orARRAY_APPEND; applies only toSubDocumentSinkHandler.
couchbase.subdocument.create.path
Controls whether a sub-document mutation creates missing parent paths.
- Type:
boolean - Default:
true - Importance: Medium
- Valid values / notes: Can be
trueorfalse; whencouchbase.subdocument.create.document=true, setting this tofalsedoes not change document creation behavior.
couchbase.subdocument.create.document
Controls whether a sub-document mutation creates a target document that does not exist.
- Type:
boolean - Default:
true - Importance: Medium
- Valid values / notes:
trueuses UPSERT storage semantics, andfalseuses REPLACE storage semantics.
couchbase.create.document
Provides a legacy compatibility alias for document creation behavior for the N1QL or sub-document handlers.
- Type:
boolean - Default: No fixed default
- Importance: Medium
- Valid values / notes: Copied to the handler setting only when the corresponding current setting is not configured.
- Deprecation status: Deprecated
- Replacement:
couchbase.n1ql.create.documentorcouchbase.subdocument.create.document
Logging and Metrics
couchbase.log.redaction
Sets the redaction level for Couchbase log markers.
- Type:
string - Default:
NONE - Importance: Medium
- Valid values / notes: Can be
NONE,PARTIAL, orFULL.
couchbase.log.document.lifecycle
Controls whether document lifecycle messages are promoted to the INFO level.
- Type:
boolean - Default:
false - Importance: Medium
- Valid values / notes: Can be
trueorfalse; enabling it may produce a high volume of logs.
couchbase.metrics.interval
Sets the interval at which the Connector writes metric snapshots to the log.
- Type:
string - Default:
10m - Importance: Medium
- Valid values / notes: Use a non-negative integer followed by
ms,s,m,h, ord;0disables metric logging. This setting has not been committed as a stable interface.
SDK Extensions
couchbase.env.*
Passes Couchbase Java SDK properties to ClusterEnvironment.
- Type:
SDK-dependent - Default:
SDK default - Importance: Medium
- Valid values / notes: After
couchbase.env., enter the SDK property name with thecom.prefix removed; an invalid property causes Task startup to fail. This pattern is not a fixed configuration item in the Connector ConfigDef.
Best Practices
Update the Same Document by a Stable Business Key
Applicable business scenario: When first connecting a continuously changing business entity, the message value already contains a stable and unique business ID. Repeated deliveries and later updates must continue to affect the same Couchbase document instead of creating multiple documents from Kafka record coordinates. Configuration example:${/orderId} extracts the document ID from the scalar field in the message value. The example also removes the field from the stored body; if the business still needs to query that field, retain the default false. If the field is missing, null, or not scalar, the Connector falls back to the Kafka key or record coordinates, so ensure before onboarding that the business ID is complete and stable.
Route Multiple Topics to Different Collections
Applicable business scenario: After connecting a single data stream, one Connector needs to consume multiple business topics and continuously write each business domain to its own Couchbase Collection while reusing the connection and authentication settings. Configuration example:\u005B and \u005D as [ and ], respectively, so the example actually uses couchbase.default.collection[orders] and couchbase.default.collection[payments]. The per-topic context configuration overrides the default Collection without a suffix; the target Scope and Collection must be created in advance. This routing is based only on the original Kafka topic and does not select a Collection dynamically from a partition or message field.
Set Limited Retries and a Durability Level for Production Writes
Applicable business scenario: The Connector already writes reliably, but transient network or target-side failures should not cause immediate failure, and Couchbase should treat a write as complete only after the specified replica confirmation condition is reached. Configuration example:2m is an example retry window; adjust it based on the recovery objective and alert response time. The Connector does not distinguish transient errors from permanent data errors, so an overly long window can delay exposure of bad records. MAJORITY requires support from the target cluster and enough available replicas. Retries may re-execute a partially successful batch, and this configuration does not provide exactly-once guarantees between Kafka offset commits and Couchbase writes; when using array appends, conditional updates, or custom non-idempotent handlers, evaluate duplicate side effects separately.
Monitoring
What to Monitor
Monitor Kafka Connect health, Connector and Task status, throughput, latency, offset commits, errors, retries, and Worker JVM heap, GC, thread, and CPU signals; monitor DLQ activity only when the corresponding error handling is enabled in the deployment.Import the Grafana Dashboard
Confirm that Kafka Connect metrics are connected to a Grafana data source and that the collected labels meet the dashboard filtering requirements; download the Kafka Connect Dashboard, import the JSON in Grafana, and select the corresponding data source.Limitations
- Couchbase writes and Kafka offset commits do not form an atomic transaction; if a failure occurs after the write succeeds but before the offset is committed, the record may be replayed after recovery, so end-to-end exactly-once guarantees are not provided.
- Different Tasks do not coordinate write ordering; records from different Kafka partitions that map to the same document ID can be written concurrently, so no global ordering across partitions, Tasks, or documents is guaranteed.
- The default handler performs full-document upserts rather than field merges; a later record with the same document ID replaces the target document with the new complete JSON.
- Sub-document array append and prepend operations are not idempotent under replay; a partially successful batch retry may insert the same element again.
- Kafka Connect per-record error tolerance and DLQ handling do not automatically handle Couchbase Handler or write errors that occur inside
SinkTask.put. - When
couchbase.durabilityis notNONE, bothcouchbase.persist.toandcouchbase.replicate.tomust beNONE.
FAQ
Why was the same business entity written as multiple documents?
Whencouchbase.document.id does not successfully extract the business ID and the Kafka record key is empty or has an unsupported type, the Connector generates a fallback ID from the topic, partition, and offset. Check that the message always contains a stable scalar business field or a stable Kafka key; when extracting from the message value, configure for example couchbase.document.id=${/orderId} and confirm that the field is not missing, null, an object, or an array.
Why was the record written to the wrong Collection?
Check thatcouchbase.default.collection uses the scope.collection or bucket.scope.collection format, and confirm that the topic name in the per-topic override exactly matches the actual record topic. Prefer couchbase.default.collection[<kafka-topic>] and avoid mixing it with the deprecated couchbase.topic.to.collection; routing does not switch Collections automatically based on a message field.
Why can the Task not connect to Couchbase after TLS is enabled?
Checkcouchbase.enable.tls, the certificate hostname, and the certificate file paths on the Worker. For a PEM CA, use couchbase.trust.certificate.path; for a Java truststore, use couchbase.trust.store.path and its password. Mutual TLS also requires the client certificate path and password. Do not disable couchbase.enable.hostname.verification unless you clearly understand the risk.
Why does a Couchbase write error still fail the Task after Kafka Connect error tolerance or a DLQ is configured?
Kafka Connect per-record error handling mainly covers stages such as the Converter and SMT. The Couchbase Handler and target write occur insideSinkTask.put, so they are not automatically skipped by errors.tolerance or written to a DLQ. Check the specific write error in the Task log and fix invalid data, permissions, the Keyspace, or the connection problem; for recoverable failures, set a limited couchbase.retry.timeout and use alerts to avoid retrying permanent errors for too long.
Why did increasing tasks.max not improve throughput, or change the update order for the same document?
Effective Task parallelism is limited by the number of Kafka topic partitions and Worker assignment; a single partition is not split among multiple Tasks. If records from different partitions use the same Couchbase document ID, different Tasks can write that document concurrently, and the final result depends on the order in which the target completes the writes. To preserve per-entity ordering, place the same entity consistently in the same Kafka partition and adjust tasks.max gradually based on target capacity.
Why did a sub-document write fail or unexpectedly create a document?
Confirm thatcouchbase.sink.handler selects SubDocumentSinkHandler and that a non-empty couchbase.subdocument.path is provided. Array append or prepend operations require the target path to satisfy Couchbase array-operation conditions; to modify existing documents only, set couchbase.subdocument.create.document to false. couchbase.subdocument.create.path controls only missing parent paths and cannot replace control over whether the target document exists.