Distributor Chart Values Reference
This reference gives an annotated example for every section of the Axual Distributor chart’s values.yaml: the initialisation job, the Kafka Connect deployment, the local and remote cluster definitions, and each of the four distribution connectors.
Type |
Reference |
Goal |
Look up what an Axual Distributor value does and what shape it takes, while writing a |
Audience |
Platform Operator writing or reviewing an Axual Distributor |
When to use |
While configuring an Axual Distributor deployment, alongside the procedure that applies these values. |
The examples below cover chart version 5.6.0. Every key the chart accepts, with its own defaults, is listed in Distributor Helm Readme, which is generated from the chart itself. The examples here show those keys in context, which the generated list cannot do.
The chart splits into three top-level areas. init creates the topics and Access Control List (ACL) entries the rest needs, connect configures the Kafka Connect deployment, and distribution configures the connectors that do the distributing.
For the procedure that applies these values, see How to Deploy the Axual Distributor.
init: the initialisation job
Kafka Connect and the distribution connector plugins both need several topics and consumer groups to exist before they start. The init section says which to create and which principal gets access to them. It reads the connect and distribution sections to work out the cluster, topic and consumer group names, so it has to agree with them.
The example uses the principal CN=Example distributor client for test instance cluster 1,OU=Infra,O=Axual. The Secret cluster-1-root holds the key data for connecting to the Apache Kafka cluster, and cluster-1-ca holds the trusted certificate authorities. The certificate in cluster-1-root needs only Alter on the cluster and Create on the timestamps topic, see Permissions the init job needs.
When the cluster uses the Principal Chain Builder, give the full chain principal instead, as in [0] CN=RootCA,OU=Infra,O=Axual, [1] CN=Intermediate CA 1,OU=Infra,O=Axual, [2] CN=Example distributor client for test instance cluster 1,OU=Infra,O=Axual. Both forms appear in the example.
|
# Initialise Kafka Cluster for Connect or Connector
init:
# Initialisation calls are executed when enabled and the connect and/or distribution is also enabled
# Default value is false
enabled: true
# Set the principal names matching the Distinguished Name of the Client Certificate
# Set multiple names if distributors connecting from other clusters to this cluster use different certificates
principals:
# Names without using Principal Chain Builder
# The principal used by a distributor connecting from cluster 1
- "User:CN=Example distributor client for test instance cluster 1,OU=Infra,O=Axual"
# The principal used by a distributor connecting from cluster 2
- "User:CN=Example distributor client for test instance cluster 2,OU=Infra,O=Axual"
# Names using Principal Chain Builder
- "User:[0] CN=RootCA,OU=Infra,O=Axual, [1] CN=Intermediate CA 1,OU=Infra,O=Axual, [2] CN=Example distributor client for test instance cluster 1,OU=Infra,O=Axual"
- "User:[0] CN=RootCA,OU=Infra,O=Axual, [1] CN=Intermediate CA 2,OU=Infra,O=Axual, [2] CN=Example distributor client for test instance cluster 2,OU=Infra,O=Axual"
# Define optional extra configuration options for the admin client. For properties, see https://kafka.apache.org/documentation/#adminclientconfigs
additionalClientConfigs:
# Set the initial TLS protocol to TLSv1.2
ssl.protocol: TLSv1.2
# Don't support pushing client metrics to server for initialisation calls
enable.metrics.push: false
# Secrets for connecting to Kafka Cluster to create topics and access control entries
tls:
# Name of and existing Secret containing the private key and certificate to access Kafka cluster
keypairSecretName: "cluster-1-root"
# The key of the Secret field containing the private key data
keypairSecretKeyName: "key"
# The key of the Secret field containing the public certificate chain for the key
keypairSecretCertName: "cert"
# Name of and existing Secret containing the list of trusted certificate authorities
truststoreCaSecretName: "cluster-1-ca"
# The key of the Secret container the trusted certificate authorities
truststoreCaSecretCertName: "cacert"
# Resource allocation for the init job pod
#resources: {}
# Defines the security options the container should be run with
#securityContext: {}
# Initialisation image information
#image:
# registry: "registry.axual.io"
# repository: "axual/distributor"
# tag: "5.6.0-0.51.0"
# pullSecrets:
# - axual-credentials
connect: the Kafka Connect deployment
This section is mostly the Strimzi KafkaConnect resource definition, so Strimzi’s own documentation applies to most of it. Five things have to be settled before it can be written:
- Kafka cluster bootstrap servers
-
The host and port pairs used for the initial connection to the local cluster. Firewall and network policies must allow Kafka Connect to reach every broker, not only the bootstrap.
- Kafka cluster trusted certificate authorities
-
Which authorities Kafka Connect trusts when connecting to the local cluster, given either as a Kubernetes Secret or as PEM data in
values.yaml. - Client certificate for the local cluster
-
The private key and certificate chain for that connection, again as a Secret or as PEM data.
- Topic and consumer group names
-
Kafka Connect needs three topics, for configuration, status and connector offsets, plus a consumer group so workers find each other. These names must be unique per Axual Distributor installation, and must not match the local cluster’s naming patterns, or the Axual Distributor would distribute its own internal data. Names conventionally begin with
_and contain the{instance-sn}short name. - Strimzi Operator version
-
The Axual Distributor’s base image has to match the operator version to avoid compatibility problems. One Axual Distributor release supports several Strimzi versions, listed in Distributor Helm Readme.
The example connects to cluster cluster01 in namespace kafka of the same Kubernetes cluster, for tenant demo and instance prod. It uses the Secrets demo-prod-distributor-client and kafka-cluster01-ca, the consumer group _demo-prod-distributor, and the Connect topics _demo-prod-distributor-config, _demo-prod-distributor-status and _demo-prod-distributor-offset. Setting strimzi.version to 0.51.0 is enough for the chart to pick the matching image tag 5.6.0-0.51.0.
# Configure the Strimzi Cluster Operator version installed
strimzi:
version: "0.51.0"
connect:
# Which image and registry should be used for connect
image:
# Optionally override the image tag (normally determined from strimzi.version)
# tag: "5.6.0-0.51.0"
repository: "axual/distributor"
# The number of connect workers to start for parallel processing
replicas: 3
# Rack aware assignment should be used
rack:
enabled: true
topologyKey: topology.kubernetes.io/zone
# The bootstrap url used to connect to the kafka cluster
bootstrapServers: "cluster01-kafka-bootstrap.kafka:9093"
groupId: "_demo-prod-distributor"
# Contains the internal connect topics settings
topics:
# Make sure to set this to a proper value for your Kafka cluster. In production this is at least 3
replicationFactor: 3
config:
name: "_demo-prod-distributor-config"
status:
name: "_demo-prod-distributor-status"
partitions: 2
offset:
name: "_demo-prod-distributor-offset"
partitions: 3
# Set extra connect configuration values
config:
key.converter: JsonConverter
value.converter: JsonConverter
# SASL PLAIN/SCRAM-SHA-25
sasl:
# Set to true to enable a SASL connection
enabled: false
type: "PLAIN"
username: "hello"
password: "world"
tls:
enabled: true
createCaCertsSecret: false
# if createTruststoreCaSecret is true, set the CA certs below
# caCerts:
# one_ca.crt: |
# -----BEGIN CERTIFICATE-----
# other_ca.crt: |
# -----BEGIN CERTIFICATE-----
# if createTruststoreCaSecret is false, the caCerts need to be set
# with an existing secret (name) and the name of the cert inside the
# secret
# caSecret:
# secretName: your_custom_ca_secret
# keyForCertificate: your_custom_ca_cert
caSecret:
secretName: "cluster01-cluster-ca-cert"
keyForCertificate: "ca.crt"
# Configure authentication using a client certificate
authentication:
enabled: true
createTlsClientSecret: false
# if createTlsClientSecret is true, set the client key and certificate chain below
#clientCert: |
# -----BEGIN CERTIFICATE-----
# -----END CERTIFICATE-----
#clientKey: |
# -----BEGIN PRIVATE KEY-----
# -----END PRIVATE KEY-----
# if a TLS secret already exists with the client credentials, provide the name here
secretName: "demo-prod-cluster01-distributor"
# Configure logging and log levels, see the Strimzi API documentation on how to configure https://strimzi.io/docs/operators/latest/configuring#type-KafkaConnectSpec-schema-reference
# If logging is not set, a default version will be deduced from the value of strimzi.version
# Log4J2 format is required for Strimzi 0.46.1 and newer
logging:
type: inline
loggers:
rootLogger.level: "INFO"
logger.distributorCommon.name: "io.axual.distributor.common"
logger.distributorCommon.level: "INFO"
logger.distributorMessage.name: "io.axual.distributor.message"
logger.distributorMessage.level: "INFO"
logger.distributorOffset.name: "io.axual.distributor.offset"
logger.distributorOffset.level: "INFO"
logger.distributorSchema.name: "io.axual.distributor.schema"
logger.distributorSchema.level: "INFO"
logger.kafkaConnectRest.name: "org.apache.kafka.connect.runtime.rest"
logger.kafkaConnectRest.level: "WARN"
# expose metrics using the JMX Prometheus Exporter
metrics:
enabled: false
# Instruct Prometheus to scrape the metrics using a PodMonitor resource
podMonitor:
enabled: false
# additionalLabels:
# label/my1: hi
# lbl2: there
# additionalAnnotations:
# annotation/my1: hi
# ann2: there
# interval: 60s
# scrapeTimeout: 12s
# Create Prometheus Alert Manager Rules
prometheusRule:
enabled: false
# additionalAnnotations:
# my-annotation: help
# some-annotation: none
# additionalLabels:
# my-label: lab-help
# some-label: 123
# rules:
# - alert: MyCustomRuleName
# annotations:
# message: '{{ "{{ $labels.connector }}" }} send rate has dropped to 0'
# expr: sum by (connector) ( kafka_connect_sink_task_sink_record_send_rate{connector=~".*-message-distributor-.*"}) == 0
# for: 5m
# labels:
# severity: high
# callout: false
# Add additional configuration here, see the Strimzi API documentation on how to use them https://strimzi.io/docs/operators/latest/configuring#type-KafkaConnectSpec-schema-reference
#resources: {}
#livenessProbe: {}
#readinessProbe: {}
#jvmOptions: {}
#jmxOptions: {}
#tracing: {}
#template: {}
#build: {}
distribution: local and remote clusters
sourceCluster describes the cluster this Axual Distributor reads from, and clusters describes every cluster it can write to. Five things have to be settled first:
- Tenant and instance short names
-
Each deployment serves one instance. The short names appear in the cluster’s naming patterns and identify the connector.
- Distribution model
-
Which clusters are the static targets at each level. See Distribution Between Kafka Clusters for the model itself.
- Target cluster connectivity
-
Bootstrap servers and security protocol per remote cluster, plus any settings needed to work around firewalls or routers that drop idle connections.
- Trusted certificate authorities per remote cluster
-
Which authorities to trust per remote cluster. Unlike the local cluster, this must be PEM data in
values.yaml. - Client certificate per remote cluster
-
The private key and certificate chain per remote cluster, also as PEM data in
values.yaml. To use Kubernetes Secrets instead, see How to Load Remote Cluster TLS Material from Kubernetes Secrets.
| Cluster names in the model are permanent. Renaming a cluster, or removing one and adding it back under a different name, can leave existing data undistributed or cause it to be distributed a second time. |
clusters is a map, so a potential target can be defined alongside the active ones. Failing over to it then means changing one reference rather than writing a new cluster definition.
|
The example has three clusters: cluster01 is the local one, cluster02 and cluster03 are remote. The Axual Distributors on cluster02 and cluster03 are not shown, and they also have to be running and distributing back to cluster01.
With one Apicurio Registry shared between clusters, set the same topicPattern on every cluster. A shared registry does not rename schemas per cluster, so differing patterns make schema lookups fail. See Schema Distribution with Apicurio.
|
distribution:
# If set to false, all distribution are disabled. If true, the individual distribution enabling flags are evaluated
enabled: true
# Setting the tenant information
tenantShortName: "demo"
instanceShortName: "prod"
# Prefix values for topic and group ACLs on the local Kafka Cluster.
prefixAclSettings:
# The local topic pattern is {tenant}-{instance}-{environment}-{topic}. environment and topic are dynamic and are ignored for the prefix ACL. Results in all topics with a name starting with this prefix to be distributed
topicPrefix: "demo-prod-"
# The local group pattern is {tenant}-{instance}-{environment}-{group.id}. environment and group id are dynamic and are ignored for the prefix ACL. Results in all offsets for a consumer group with a name starting with this prefix to be distributed
groupPrefix: "demo-prod-"
# Information about the source cluster. This is the cluster Kafka Connect is connecting with and distributors read their data from
sourceCluster:
# Cluster name, this is also used in the connector and consumer group names
name: cluster01
# Contains instance topic information for the cluster
topics:
# The replication factor to use for the local distributor specific topics
# Make sure to set this to a proper value for your Kafka cluster. In production this is at least 3
replicationFactor: 3
# Timestamps topic for offset distribution, name is _<tenant>-<instance>-consumer-timestamps
timestamps:
nameOverride: ""
# The number of partitions to use for the timestamps topic
partitions: 25
# The resource naming pattern used on this cluster
topicPattern: '{tenant}-{instance}-{environment}-{topic}'
groupPattern: '{tenant}-{instance}-{environment}-{group.id}'
# If a pattern value is missing the value specified in this dictionary should be used
patternDefaultValues:
environment: prod
# Connection information, used by offset distributor and committer to update/query data from the source cluster
bootstrapServers: "cluster01-kafka-bootstrap.kafka:9093"
# What is the protocol for the Kafka connection using the bootstrap server.
# Valid values are PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL
securityProtocol: "SSL"
# Authentication configuration is provided separately, but other client connection settings
# can be specified here. For example, TLS protocols and buffer sizes
additionalClientConfigs:
ssl.enabled.protocols: TLSv1.2
# If the security protocol is SASL_PLAINTEXT or SASL_SSL then this section must be used
# sasl:
# Which SASL mechanism should be used. Accepts PLAIN, SCRAM-SHA-256 or SCRAM-SHA-512
# mechanism: SCRAM-SHA-512
# username: my-user
# password: my-password
# Contains the tls settings for the cluster
tls:
# A list of PEM encoded CA certificates to use when connecting to the Kafka cluster
# This is used to verify the server certificate
caCerts:
- |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# The PEM encoded client certificate chain to use if Mutual TLS, or certificate based authentication is required
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
clientCert: |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# A PEM encoded PKCS#8 key to use if Mutual TLS, or certificate based authentication is required
# Encrypted keys are supported, the password should go in the clientKeyPassword configuration
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
clientKey: |
-----BEGIN PRIVATE KEY-----
-----END PRIVATE KEY-----
# The password for the clientKey, only set this if the clientKey is encrypted
# clientKeyPassword: somethingVerySecret
# Information about all potential target clusters. These are the cluster where data will be written to
clusters:
# The key/id of the cluster to use. Recommend to keep this the same as the cluster name
cluster02:
# Cluster name, this is also used in the connector and consumer group names
name: cluster02
# Contains instance topic information for the cluster
topics:
# Timestamps topic for offset distribution
timestamps:
# Use name override if the target cluster does not use default name of _<tenant>-<instance>-consumer-timestamps
nameOverride: ""
# Schema topic for schema distribution
schemas:
# Use name override if the target cluster does not use default name of _<tenant>-<instance>-schemas
nameOverride: ""
# The resource naming pattern used on this cluster
topicPattern: '{tenant}-{instance}-{environment}-{topic}'
groupPattern: '{tenant}-{instance}-{environment}-{group.id}'
# If a pattern value is missing the value specified in this dictionary should be used
patternDefaultValues:
environment: prod
# Connection information, used by connectors to connect to the target
bootstrapServers: "cluster02-kafka-bootstrap.region2.demo.com:9093"
# What is the protocol for the Kafka connection using the bootstrap server.
# Valid values are PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL
securityProtocol: SSL
# Authentication configuration is provided separately, but other client connection settings
# can be specified here. For example, TLS protocols and buffer sizes
additionalClientConfigs:
ssl.enabled.protocols: TLSv1.2
# If the security protocol is SASL_PLAINTEXT or SASL_SSL then this section must be used
#sasl:
# # Which SASL mechanism should be used. Accepts PLAIN, SCRAM-SHA-256 or SCRAM-SHA-512
# mechanism: SCRAM-SHA-512
# username: my-user
# password: my-password
# Contains the tls settings for the cluster
tls:
# A list of PEM encoded CA certificates to use when connecting to the target Kafka cluster
# This is used to verify the server certificate
caCerts:
- |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# The PEM encoded client certificate chain to use if Mutual TLS, or certificate based authentication is required
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
clientCert: |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# A PEM encoded PKCS#8 key to use if Mutual TLS, or certificate based authentication is required
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
# Encrypted keys are supported, the password should go in the clientKeyPassword configuration
clientKey: |
-----BEGIN PRIVATE KEY-----
-----END PRIVATE KEY-----
# The password for the clientKey, only set this if the clientKey is encrypted
#clientKeyPassword: somethingVerySecret
# The key/id of the cluster to use. Recommend to keep this the same as the cluster name
cluster03:
# Cluster name, this is also used in the connector and consumer group names
name: cluster03
# Contains instance topic information for the cluster
topics:
# Timestamps topic for offset distribution
timestamps:
# Use name override to not use default name of _<tenant>-<instance>-consumer-timestamps
nameOverride: ""
# Schema topic for schema distribution
schemas:
# Use name override to not use default name of _<tenant>-<instance>-schemas
nameOverride: ""
# The resource naming pattern used on this cluster
topicPattern: '{tenant}-{instance}-{environment}-{topic}'
groupPattern: '{tenant}-{instance}-{environment}-{group.id}'
# If a pattern value is missing the value specified in this dictionary should be used
patternDefaultValues:
environment: prod
# Connection information, used by connectors to connect to the target
bootstrapServers: "cluster03-kafka-bootstrap.region3.demo.com:9093"
# What is the protocol for the Kafka connection using the bootstrap server.
# Valid values are PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL
securityProtocol: SSL
# Authentication configuration is provided separately, but other client connection settings
# can be specified here. For example, TLS protocols and buffer sizes
additionalClientConfigs:
ssl.enabled.protocols: TLSv1.2
# If the security protocol is SASL_PLAINTEXT or SASL_SSL then this section must be used
#sasl:
# # Which SASL mechanism should be used. Accepts PLAIN, SCRAM-SHA-256 or SCRAM-SHA-512
# mechanism: SCRAM-SHA-512
# username: my-user
# password: my-password
# Contains the tls settings for the cluster
tls:
# A list of PEM encoded CA certificates to use when connecting to the target Kafka cluster
# This is used to verify the server certificate
caCerts:
- |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# The PEM encoded client certificate chain to use if Mutual TLS, or certificate based authentication is required
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
clientCert: |
-----BEGIN CERTIFICATE-----
-----END CERTIFICATE-----
# A PEM encoded PKCS#8 key to use if Mutual TLS, or certificate based authentication is required
# Only when both clientCert and clientKey are provided will Mutual TLS be enabled
# Encrypted keys are supported, the password should go in the clientKeyPassword configuration
clientKey: |
-----BEGIN PRIVATE KEY-----
-----END PRIVATE KEY-----
# The password for the clientKey, only set this if the clientKey is encrypted
#clientKeyPassword: somethingVerySecret
schemaDistributor
Reads the schema topic on the local cluster and produces it to the clusters named in targetClusterNames, taking their connection details and topic names from the sourceCluster and clusters sections above. additionalClientConfigs here is merged with the target cluster’s own additionalClientConfigs to configure the producer.
| The Schema Distributor works only with the Confluent-based Legacy Schema Registry. With Apicurio Registry, leave it disabled and share one Apicurio instance between clusters instead. |
| Enable the Schema Distributor on the primary cluster only. It synchronises from a leading registry to following registries, so running it on more than one cluster can loop, delete schemas from the leading topic, and produce mismatched schema ids. |
distribution:
# Schema Distribution
schemaDistributor:
enabled: true
# Which clusters should it distribute to, uses the {clusters.[].name} to find cluster information
targetClusterNames:
- "cluster02"
- "cluster03"
# Override the settings for the internal Kafka Connect consumer for the tasks. This can be used to tweak performance
consumerOverrides:
# Override how many records will be read in a single poll of Kafka Connect
max.poll.records: 50
# Any additional client configs needed to connect to the remote cluster.
# For example, TLS protocols and buffer sizes and retries
additionalClientConfigs:
acks: all
retries: 5
max.in.flight.requests.per.connection: 5
enable.idempotence: true
ssl.enabled.protocols: TLSv1.2,TLSv1.3
ssl.endpoint.identification.algorithm: ''
offsetCommitter
Reads offset timestamps for consumer groups from the local cluster’s timestamps topic, translates them to offsets, and commits those offsets locally. It can be enabled before anything else, because it consumes what the remote clusters' Offset Distributors send it and waits while none are running.
distributionOffsetMs is subtracted from each received timestamp before the offset is resolved, which is what stops a client missing records when it migrates between clusters. Raise it when distribution to this cluster is slow. Resolving group and topic names is expensive, so topicCacheSize and groupCacheSize bound the caches that hold the resolved names. additionalClientConfigs here is merged with sourceCluster.additionalClientConfigs.
distribution:
offsetCommitter:
enabled: true
# how many parallel tasks should the connector have
maxTasks: 5
# Setting this to true will skip any connectivity check when creating and validating the connector config
# skipConnectionValidation: true
# On severe failures there will be continuous processing errors.
# Internal connector resources can be reinitialized when reaching the thresholds defined here
reinitializeOnContinuousError:
# Reinitialize if the continuous error count reaches this number
maximumCount: 5000
# Reinitialize if the continues error count has existed for this long
maximumTimeMs: 30000
# The offset committer reads timestamps that relate to records produced on the source topic.
# This value pushed the timestamp back to make sure that any application will continue before or at
# the offset of the original record, guaranteeing at-least-once consuming of records
# If the message distribution to this cluster slows down a higher value might be needed
distributionOffsetMs: 30000
# The maximum size of the cache used to optimize group and topic resolving.
topicCacheSize: 1000
groupCacheSize: 1000
# Override the settings for the internal Kafka Connect consumer for the tasks. This can be used to tweak performance
consumerOverrides:
# Override how many records will be read in a single poll of Kafka Connect
max.poll.records: 50
# Any additional client configs needed to connect to the source cluster to commit the offsets.
# For example, TLS protocols and buffer sizes
additionalClientConfigs:
ssl.enabled.protocols: TLSv1.2,TLSv1.3
ssl.endpoint.identification.algorithm: ''
levels: message and offset distributors
The Message Distributor decides per record whether to distribute it, using the distribution model. Its configuration is per level, numbered 1 to 32, and each level names one target cluster and carries both a Message Distributor and an Offset Distributor block.
The Message Distributor needs topicRegex to subscribe to the local topics it should distribute, and that expression has to match the cluster’s topic pattern or distribution breaks. In the example the pattern is {tenant}-{instance}-{environment}-{topic} with demo and prod fixed, so demo-prod-.* matches every topic in the instance. Level 1 distributes to cluster02 with 5 tasks and level 2 to cluster03 with 2 tasks.
levels.[].messageDistributor.additionalTargetClientConfigs is merged with the target cluster’s additionalClientConfigs to configure the producer to the remote cluster.
distribution:
# Axual uses a multi level distribution model to make sure records and offsets arrive at the correct cluster
# This is a dictionary containing the distribution levels relevant to the source cluster.
# Each level can be activated independent and contains settings for the offset and message distributors for that
# level
# The level number is 1 to 32
levels:
# level numbers should be numbers, and are used in connector and group naming
1:
enabled: true
# To which cluster does the data go, uses the {clusters.[].name} to find cluster information
targetClusterName: "cluster02"
# The message distributor settings to copy records from the source cluster to the target cluster
messageDistributor:
enabled: true
# how many parallel tasks should the connector have
maxTasks: 5
# which topic pattern should be used for subscription. Form is dependent on to {sourceCluster.topicPattern}
topicRegex: 'demo-prod-.*'
# The maximum size of the cache used to optimize topic resolving.
topicCacheSize: 1000
# Override the settings for the internal Kafka Connect consumer for the tasks. This can be used to tweak performance
consumerOverrides:
# Override how many records will be read in a single poll of Kafka Connect
max.poll.records: 50
# Any additional client configs needed to connect to the target cluster.
# For example, TLS protocols and buffer sizes
additionalTargetClientConfigs: {}
2:
enabled: true
targetClusterName: "cluster03"
messageDistributor:
enabled: true
maxTasks: 2
topicRegex: 'demo-prod-.*'
topicCacheSize: 1000
consumerOverrides:
max.poll.records: 50
additionalTargetClientConfigs: {}
Configuring Offset Distributor
The Offset Distributor uses the distribution model to decide whether a committed offset should be distributed. Like the Message Distributor, it is configured per level, and each level holds the settings for one remote cluster.
It reads __consumer_offsets on the local cluster and drops every commit that belongs to another tenant or instance, or that the distribution model excludes. A consumer group can commit often, because the consumer application and the partition load decide that, so the Offset Distributor produces at most one timestamp per consumer group per time window. windowSizeMs sets that window.
Two configuration maps merge with a cluster-level map, one per side of the connection:
| Key | Merged with | Configures |
|---|---|---|
|
|
The clients that look up and commit offsets on the local cluster. |
|
|
The client that produces timestamps to the remote cluster. |
The example configures level 1 distribution to cluster02 with 5 tasks and level 2 distribution to cluster03 with 5 tasks. The window size for both is 15 seconds (15000 milliseconds).
distribution:
# Axual uses a multi level distribution model to make sure records and offsets arrive at the correct cluster
# This is a dictionary containing the distribution levels relevant to the source cluster.
# Each level can be activated independent and contains settings for the offset and message distributors for that
# level
# The level number is 1 to 32
levels:
# level numbers should be numbers, and are used in connector and group naming
1:
enabled: true
# To which cluster does the data go, uses the {clusters.[].name} to find cluster information
targetClusterName: "cluster02"
# The offset distributor settings to distribute offsets from the source cluster as a timestamp
# to the target cluster where it can be committed.
# The distributor produces to the timestamps topic on te target cluster
# To prevent a flood of offset commits only the latest value for a specific consumer group and topic partition
# combination in a time window is translated and produced as a timestamp to the remote cluster
offsetDistributor:
enabled: true
# how many parallel tasks should the connector have
maxTasks: 5
# The size of the window to group offset commits in.
# Grouping is done on group.id+topic+partition.
windowSizeMs: 15000
# The maximum size of the cache used to optimize group and topic resolving.
topicCacheSize: 1000
groupCacheSize: 1000
# Override the settings for the internal Kafka Connect consumer for the tasks. This can be used to tweak performance
consumerOverrides:
# Override how many records will be read in a single poll of Kafka Connect
max.poll.records: 50
# Any additional client configs needed to connect to the source cluster.
# For example, TLS protocols and buffer sizes
# The local connection is used to determine the timestamps of a specific record
additionalLocalClientConfigs: {}
# Any additional client configs needed to connect to the target cluster.
# For example, TLS protocols and buffer sizes
additionalTargetClientConfigs: {}
2:
enabled: true
targetClusterName: "cluster03"
offsetDistributor:
enabled: true
maxTasks: 5
windowSizeMs: 15000
topicCacheSize: 1000
groupCacheSize: 1000
consumerOverrides:
max.poll.records: 50
additionalLocalClientConfigs: {}
additionalTargetClientConfigs: {}