Go Dcp Elasticsearch

Go implementation of
the Elasticsearch Connect Couchbase.
Go Dcp Elasticsearch streams documents from Couchbase Database Change Protocol (DCP) and writes to
Elasticsearch index in near real-time. You can find more information by looking at our docs.
Features
- Less resource usage and higher throughput(see Benchmarks).
- Custom routing support(see Example).
- Update multiple documents for a DCP event(see Example).
- Handling different DCP events such as expiration, deletion and mutation(see Example).
- Elasticsearch compression request body support.
- Managing batch configurations such as maximum batch size, batch bytes, batch ticker durations.
- Scale up and down by custom membership algorithms(Couchbase, KubernetesHa, Kubernetes StatefulSet or
Static, see examples).
- Multiple Elasticsearch clusters — route actions via
ClusterKey (see Multiple Elasticsearch clusters).
- Easily manageable configurations.
Benchmarks
The benchmark was made with the 1,001,006 Couchbase document, because it is possible to more clearly observe the
difference in the batch structure between the two packages. Default configurations for Java Elasticsearch Connect
Couchbase
used for both connectors.
| Package |
Time to Process Events |
Elasticsearch Indexing Rate(/s) |
Average CPU Usage(Core) |
Average Memory Usage |
| Go Dcp Elasticsearch(Go 1.20) |
50s |
 |
0.486 |
408MB |
| Java Elasticsearch Connect Couchbase(JDK15) |
80s |
 |
0.31 |
1091MB |
Example
Struct Config
func mapper(event couchbase.Event) []document.ESActionDocument {
if event.IsMutated {
e := document.NewIndexAction(event.Key, event.Value, nil)
return []document.ESActionDocument{e}
}
e := document.NewDeleteAction(event.Key, nil)
return []document.ESActionDocument{e}
}
func main() {
connector, err := dcpelasticsearch.NewConnectorBuilder(config.Config{
Elasticsearch: config.Elasticsearch{
CollectionIndexMapping: map[string]string{
"_default": "indexname",
},
Urls: []string{"http://localhost:9200"},
},
Dcp: dcpConfig.Dcp{
Username: "user",
Password: "password",
BucketName: "dcp-test",
Hosts: []string{"localhost:8091"},
Dcp: dcpConfig.ExternalDcp{
Group: dcpConfig.DCPGroup{
Name: "groupName",
},
},
Metadata: dcpConfig.Metadata{
Config: map[string]string{
"bucket": "checkpoint-bucket-name",
"scope": "_default",
"collection": "_default",
},
Type: "couchbase",
},
},
}).
SetMapper(mapper).
Build()
if err != nil {
panic(err)
}
defer connector.Close()
connector.Start()
}
File Config
Default Mapper
Multiple Elasticsearch clusters
Configuration
Dcp Configuration
Check out on go-dcp
Elasticsearch Specific Configuration
| Variable |
Type |
Required |
Default |
Description |
elasticsearch.collectionIndexMapping |
map[string]string |
yes |
|
Defines which Couchbase collection events will be written to which index |
elasticsearch.urls |
[]string |
yes |
|
Elasticsearch connection urls |
elasticsearch.username |
string |
no |
|
The username of Elasticsearch |
elasticsearch.password |
string |
no |
|
The password of Elasticsearch |
elasticsearch.typeName |
string |
no |
|
Defines Elasticsearch index type name |
elasticsearch.batchSizeLimit |
int |
no |
1000 |
Maximum message count for batch, if exceed flush will be triggered. |
elasticsearch.batchTickerDuration |
time.Duration |
no |
10s |
Batch is being flushed automatically at specific time intervals for long waiting messages in batch. |
elasticsearch.batchCommitTickerDuration |
time.Duration |
no |
0s |
Configures checkpoint offset save time, By default, after batch flushing, the offsets are updated immediately, this period can be increased for performance. |
elasticsearch.batchByteSizeLimit |
int, string |
no |
10mb |
Maximum size(byte) for batch, if exceed flush will be triggered. 10mb is default. |
elasticsearch.maxConnsPerHost |
int |
no |
512 |
Maximum number of connections per each host which may be established |
elasticsearch.maxIdleConnDuration |
time.Duration |
no |
10s |
Idle keep-alive connections are closed after this duration. |
elasticsearch.compressionEnabled |
boolean |
no |
false |
Compression can be used if message size is large, CPU usage may be affected. |
elasticsearch.concurrentRequest |
int |
no |
1 |
Concurrent bulk request count |
elasticsearch.disableDiscoverNodesOnStart |
boolean |
no |
false |
Disable discover nodes when initializing the client. |
elasticsearch.discoverNodesInterval |
time.Duration |
no |
5m |
Discover nodes periodically |
elasticsearch.rejectionLog.index |
string |
no |
cbes-rejects |
Rejection log index name. cbes-rejects is default. |
elasticsearch.rejectionLog.includeSource |
boolean |
no |
false |
Includes rejection log source info. false is default. |
elasticsearch.maxRetries |
int |
no |
math.MaxInt |
Maximum retry count for the Elasticsearch client (per bulk sub-request). |
elasticsearch.retry.enabled |
boolean |
no |
false |
Enables the built-in retry layer that re-submits only the retryable items of a failed bulk request. Disabled by default. |
elasticsearch.retry.maxRetries |
int |
no |
3 |
Maximum retry attempts for retryable failures before falling through to OnError/panic. |
elasticsearch.retry.retryOnStatus |
[]int |
no |
[429,502,503,504] |
HTTP status codes treated as retryable (both per-item and whole-response). Everything else is terminal. |
elasticsearch.retry.initialInterval |
time.Duration |
no |
200ms |
Starting backoff before the first retry; grows exponentially with full jitter. |
elasticsearch.retry.maxInterval |
time.Duration |
no |
5s |
Upper bound on the backoff between retries. |
elasticsearch.clusters |
map[string]object |
no |
|
Optional named Elasticsearch clusters. Each entry mirrors elasticsearch connection fields (urls, auth, collectionIndexMapping, retry, …). Use document.ESActionDocument.ClusterKey to route an action to a name defined here. |
elasticsearch.rejectionLog.targetCluster |
string |
no |
|
When using RejectionLogSinkResponseHandler, writes rejection documents via the client for this cluster key (empty = default cluster). |
elasticsearch.tls.skipVerify |
bool |
no |
|
If set to true, Elasticsearch client will skip TLS verification. Only set to true on dev environments. |
elasticsearch.tls.caCert |
[]byte |
no |
|
CA certificate bytes. |
elasticsearch.tls.cert |
[]byte |
no |
|
Client certificate bytes. |
elasticsearch.tls.key |
[]byte |
no |
|
Key file bytes. |
Multiple Elasticsearch clusters
The primary block under elasticsearch is the default cluster (empty ClusterKey). Optional elasticsearch.clusters defines named clusters with the same shape as the root block (at minimum urls; use collectionIndexMapping per cluster when resolving index names from Couchbase collections).
In your mapper, set ClusterKey on each document.ESActionDocument to a name from elasticsearch.clusters. Leave it empty to use the default cluster. The reserved name default is normalized to the primary cluster.
Example: example/multi-cluster/main.go
Retryable bulk failures
By default a failing bulk request is surfaced to the registered SinkResponseHandler, or — when none is registered — the connector panics. Enabling elasticsearch.retry adds an optional layer that first re-submits only the retryable items of a failed bulk (both transport/connection errors and configurable per-item / whole-response statuses such as 429, 503) with exponential backoff and jitter. Terminal failures (e.g. 4xx validation errors) are never retried. After maxRetries is exhausted the remaining failures fall through to OnError/panic exactly as before, so DCP replay semantics are preserved.
Retry is configured per cluster: the block under elasticsearch applies to the default cluster, and each named cluster under elasticsearch.clusters may declare its own retry block. A named cluster that omits retry inherits the default cluster's retry settings. A transient failure is retried only against the cluster that produced it, so one unhealthy cluster does not restart the whole fleet.
When retry is enabled the Elasticsearch client's own transport-level retry is disabled for that cluster, so all retry/backoff is governed by this layer. Otherwise the client would retry 5xx/transport failures itself with no backoff and a near-infinite retry count, shadowing the configured maxRetries/backoff. When retry is disabled, the client's default retry behavior is unchanged.
elasticsearch:
urls: [ "http://localhost:9200" ]
retry:
enabled: true
maxRetries: 3
retryOnStatus: [429, 502, 503, 504]
initialInterval: 200ms
maxInterval: 5s
clusters:
analytics:
urls: [ "http://localhost:9201" ]
# no retry block -> inherits the default cluster's retry settings
Exposed metrics
| Metric Name |
Description |
Labels |
Value Type |
| cbgo_elasticsearch_connector_latency_ms_current |
Time to adding to the batch. |
N/A |
Gauge |
| cbgo_elasticsearch_connector_bulk_request_process_latency_ms_current |
Time to process bulk request. |
N/A |
Gauge |
| cbgo_elasticsearch_connector_action_total_current |
Count elasticsearch actions |
action_type: Type of action (e.g., delete, index) result: Result of the action (e.g., success, error) index_name: The name of the index to which the action is applied |
Counter |
You can also use all DCP-related metrics explained here.
All DCP-related metrics are automatically injected. It means you don't need to do anything.
Grafana Metric Dashboard
Grafana & Prometheus Example
Contributing
Go Dcp Elasticsearch is always open for direct contributions. For more information please check
our Contribution Guideline document.
License
Released under the MIT License.