Kafka Share
Consume messages from Apache Kafka topics as a queue, using a share group.
What’s inside
-
Kafka Share component, URI syntax:
kafka-share:topic
Please refer to the above links for usage and configuration details.
Maven coordinates
<dependency>
<groupId>org.apache.camel.springboot</groupId>
<artifactId>camel-kafka-share-starter</artifactId>
</dependency> Spring Boot Auto-Configuration
The starter supports 80 options, which are listed below.
| Name | Description | Default | Type |
|---|---|---|---|
camel.component.kafka-share.acquire-mode | How the share consumer acquires records. With record_limit, a poll() returns at most maxPollRecords records. With batch_optimized, a poll() can return more records than maxPollRecords, to align with the batches of the topic. | batch_optimized | String |
camel.component.kafka-share.additional-properties | Sets additional properties for either kafka consumer or kafka producer in case they can’t be set directly on the camel configurations (e.g.: new Kafka properties that are not reflected yet in Camel configurations), the properties have to be prefixed with additionalProperties., e.g.: additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro. If the properties are set in the application.properties file, they must be prefixed with camel.component.kafka.additional-properties (camel.component.kafka-share.additional-properties for the Kafka Share component) followed by the property name enclosed in square brackets, for example the delivery.timeout.ms property in square brackets. This is a multi-value option with prefix: additionalProperties. | Object> | |
camel.component.kafka-share.autowired-enabled | Whether autowiring is enabled. This is used for automatic autowiring options (the option must be marked as autowired) by looking up in the registry to find if there is a single instance of matching type, which then gets configured on the component. This can be used for automatic configuring JDBC data sources, JMS connection factories, AWS Clients, etc. | true | Boolean |
camel.component.kafka-share.bridge-error-handler | Allows for bridging the consumer to the Camel routing Error Handler, which mean any exceptions (if possible) occurred while the Camel consumer is trying to pickup incoming messages, or the likes, will now be processed as a message and handled by the routing Error Handler. Important: This is only possible if the 3rd party component allows Camel to be alerted if an exception was thrown. Some components handle this internally only, and therefore bridgeErrorHandler is not possible. In other situations we may improve the Camel component to hook into the 3rd party component and make this possible for future releases. By default the consumer will use the org.apache.camel.spi.ExceptionHandler to deal with exceptions, that will be logged at WARN or ERROR level and ignored. | false | Boolean |
camel.component.kafka-share.brokers | URL of the Kafka brokers to use. The format is host1:port1,host2:port2, and the list can be a subset of brokers or a VIP pointing to a subset of brokers. This option is known as bootstrap.servers in the Kafka documentation. | String | |
camel.component.kafka-share.check-crcs | Automatically check the CRC32 of the records consumed. This ensures no on-the-wire or on-disk corruption to the messages occurred. | true | Boolean |
camel.component.kafka-share.client-id | The client id is a user-specified string sent in each request to help trace calls. It should logically identify the application making the request. | String | |
camel.component.kafka-share.commit-mode | How the acknowledgements of a poll are committed to the broker. SYNC commits them after the records of the poll are processed, and waits for the result, so a failure to commit is reported to the exception handler. ASYNC commits them without waiting, and a failure to commit is logged. | sync | KafkaShareCommitMode |
camel.component.kafka-share.commit-timeout-ms | The maximum time to wait for the acknowledgements to be committed, when commitMode is SYNC. The option is a java.lang.Long type. | 5000 | Long |
camel.component.kafka-share.configuration | Allows to pre-configure the Kafka share component with common options that the endpoints will reuse. The option is a org.apache.camel.component.kafka.share.KafkaShareConfiguration type. | KafkaShareConfiguration | |
camel.component.kafka-share.connection-max-idle-ms | Close idle connections after the number of milliseconds specified by this config. | 540000 | Integer |
camel.component.kafka-share.consumer-request-timeout-ms | The configuration controls the maximum amount of time the client will wait for the response of a request. If the response is not received before the timeout elapses, the client will resend the request if necessary or fail the request if retries are exhausted. | 30000 | Integer |
camel.component.kafka-share.consumers-count | The number of consumers that connect to the Kafka server. Each consumer runs on its own thread and receives records, as the records of a partition are shared by all the consumers of a share group. Unlike a consumer group, the number of consumers is not limited by the number of partitions. | 1 | Integer |
camel.component.kafka-share.create-consumer-backoff-interval | The delay in millis seconds to wait before trying again to create the kafka consumer (kafka-client). | 5000 | Long |
camel.component.kafka-share.create-consumer-backoff-max-attempts | Maximum attempts to create the kafka consumer (kafka-client), before eventually giving up and failing. Error during creating the consumer may be fatal due to invalid configuration and as such recovery is not possible. However, one part of the validation is DNS resolution of the bootstrap broker hostnames. This may be a temporary networking problem, and could potentially be recoverable. While other errors are fatal, such as some invalid kafka configurations. Unfortunately, kafka-client does not separate this kind of errors. Camel will by default retry forever, and therefore never give up. If you want to give up after many attempts then set this option and Camel will then when giving up terminate the consumer. To try again, you can manually restart the consumer by stopping, and starting the route. | Integer | |
camel.component.kafka-share.enabled | Whether to enable auto configuration of the kafka-share component. This is enabled by default. | Boolean | |
camel.component.kafka-share.fetch-max-bytes | The maximum amount of data the server should return for a fetch request. This is not an absolute maximum: if the first record batch in the first non-empty partition of the fetch is larger than this value, the record batch will still be returned to ensure that the consumer can make progress. | 52428800 | Integer |
camel.component.kafka-share.fetch-min-bytes | The minimum amount of data the server should return for a fetch request. If insufficient data is available, the request will wait for that much data to accumulate before answering the request. | 1 | Integer |
camel.component.kafka-share.fetch-wait-max-ms | The maximum amount of time the server will block before answering the fetch request if there isn’t enough data to immediately satisfy fetch.min.bytes. | 500 | Integer |
camel.component.kafka-share.group-id | The name of the share group. All the consumers that use the same share group name share the records of the topics: each record is delivered to one of them. | String | |
camel.component.kafka-share.header-deserializer | To use a custom KafkaHeaderDeserializer to deserialize kafka headers values. The option is a org.apache.camel.component.kafka.serde.KafkaHeaderDeserializer type. | KafkaHeaderDeserializer | |
camel.component.kafka-share.header-filter-strategy | To use a custom HeaderFilterStrategy to filter header to and from Camel message. The option is a org.apache.camel.spi.HeaderFilterStrategy type. | HeaderFilterStrategy | |
camel.component.kafka-share.health-check-consumer-enabled | Used for enabling or disabling all consumer based health checks from this component | true | Boolean |
camel.component.kafka-share.health-check-producer-enabled | Used for enabling or disabling all producer based health checks from this component. Notice: Camel has by default disabled all producer based health-checks. You can turn on producer checks globally by setting camel.health.producersEnabled=true. | true | Boolean |
camel.component.kafka-share.kafka-share-client-factory | Factory to use for creating org.apache.kafka.clients.consumer.KafkaShareConsumer instances. This allows configuring a custom factory to create instances with logic that extends the vanilla Kafka clients. The option is a org.apache.camel.component.kafka.share.KafkaShareClientFactory type. | KafkaShareClientFactory | |
camel.component.kafka-share.kerberos-before-relogin-min-time | Login thread sleep time between refresh attempts. | 60000 | Integer |
camel.component.kafka-share.kerberos-config-location | Location of the kerberos config file. | String | |
camel.component.kafka-share.kerberos-init-cmd | Kerberos kinit command path. Default is /usr/bin/kinit | /usr/bin/kinit | String |
camel.component.kafka-share.kerberos-principal-to-local-rules | A list of rules for mapping from principal names to short names (typically operating system usernames). The rules are evaluated in order, and the first rule that matches a principal name is used to map it to a short name. Any later rules in the list are ignored. By default, principal names of the form {username}/{hostname}{REALM} are mapped to {username}. For more details on the format, please see the Security Authorization and ACLs documentation (at the Apache Kafka project website). Multiple values can be separated by comma | DEFAULT | String |
camel.component.kafka-share.kerberos-renew-jitter | Percentage of random jitter added to the renewal time. | Double | |
camel.component.kafka-share.kerberos-renew-window-factor | Login thread will sleep until the specified window factor of time from last refresh to ticket’s expiry has been reached, at which time it will try to renew the ticket. | Double | |
camel.component.kafka-share.key-deserializer | Deserializer class for the key that implements the Deserializer interface. | org.apache.kafka.common.serialization.StringDeserializer | String |
camel.component.kafka-share.max-partition-fetch-bytes | The maximum amount of data per-partition the server will return. | 1048576 | Integer |
camel.component.kafka-share.max-poll-records | The maximum number of records returned in a single call to poll(). With the batch_optimized acquire mode, a poll can return more records, to align with the batches of the topic. | 500 | Integer |
camel.component.kafka-share.metadata-max-age-ms | The period of time in milliseconds after which we force a refresh of metadata even if we haven’t seen any partition leadership changes to proactively discover any new brokers or partitions. | 300000 | Integer |
camel.component.kafka-share.metric-reporters | A list of classes to use as metrics reporters. Implementing the MetricReporter interface allows plugging in classes that will be notified of new metric creation. The JmxReporter is always included to register JMX statistics. | String | |
camel.component.kafka-share.metrics-sample-window-ms | The window of time a metrics sample is computed over. | 30000 | Integer |
camel.component.kafka-share.no-of-metrics-sample | The number of samples maintained to compute metrics. | 2 | Integer |
camel.component.kafka-share.oauth-client-id | OAuth client ID. Used when saslAuthType is set to OAUTH. | String | |
camel.component.kafka-share.oauth-client-secret | OAuth client secret. Used when saslAuthType is set to OAUTH. | String | |
camel.component.kafka-share.oauth-scope | OAuth scope. Used when saslAuthType is set to OAUTH. | String | |
camel.component.kafka-share.oauth-token-endpoint-uri | OAuth token endpoint URI. Used when saslAuthType is set to OAUTH. | String | |
camel.component.kafka-share.on-failure | How to acknowledge a record whose exchange failed or was rolled back. RELEASE makes the record available again, to this or another consumer, until the broker delivery count limit is reached. REJECT discards the record. ACCEPT marks the record as consumed. | release | KafkaShareAcknowledgeType |
camel.component.kafka-share.poll-exception-strategy | To use a custom strategy with the consumer to control how to handle exceptions thrown from the Kafka broker while polling messages. The option is a org.apache.camel.component.kafka.PollExceptionStrategy type. | PollExceptionStrategy | |
camel.component.kafka-share.poll-on-error | What to do if the share consumer throws an exception while polling for new records. DISCARD and RETRY log the exception and poll again: unlike the kafka component, there is no record to skip or to retry, as the poll itself failed, and the records that fail in the route are handled with onFailure. ERROR_HANDLER lets the exception handler of the consumer handle the exception, and polls again. RECONNECT closes the share consumer and creates a new one. STOP stops consuming. An authentication or authorization failure always stops consuming. | error-handler | PollOnError |
camel.component.kafka-share.poll-timeout-ms | The timeout used when polling the share consumer. The option is a java.lang.Long type. | 5000 | Long |
camel.component.kafka-share.pre-validate-host-and-port | Whether to eager validate that broker host:port is valid and can be DNS resolved to known host during starting this consumer. If the validation fails, then an exception is thrown, which makes Camel fail fast. Disabling this will postpone the validation after the consumer is started, and Camel will keep re-connecting in case of validation or DNS resolution error. | true | Boolean |
camel.component.kafka-share.receive-buffer-bytes | The size of the TCP receive buffer (SO_RCVBUF) to use when reading data. | 65536 | Integer |
camel.component.kafka-share.reconnect-backoff-max-ms | The maximum amount of time in milliseconds to wait when reconnecting to a broker that has repeatedly failed to connect. If provided, the backoff per host will increase exponentially for each consecutive connection failure, up to this maximum. After calculating the backoff increase, 20% random jitter is added to avoid connection storms. | 1000 | Integer |
camel.component.kafka-share.reconnect-backoff-ms | The amount of time to wait before attempting to reconnect to a given host. This avoids repeatedly connecting to a host in a tight loop. This backoff applies to all requests sent by the consumer to the broker. | 50 | Integer |
camel.component.kafka-share.retry-backoff-max-ms | The maximum amount of time in milliseconds to wait when retrying a request to the broker that has repeatedly failed. If provided, the backoff per client will increase exponentially for each failed request, up to this maximum. To prevent all clients from being synchronized upon retry, a randomized jitter with a factor of 0.2 will be applied to the backoff, resulting in the backoff falling within a range between 20% below and 20% above the computed value. If retry.backoff.ms is set to be higher than retry.backoff.max.ms, then retry.backoff.max.ms will be used as a constant backoff from the beginning without any exponential increase | 1000 | Integer |
camel.component.kafka-share.retry-backoff-ms | The amount of time to wait before attempting to retry a failed request to a given topic partition. This avoids repeatedly sending requests in a tight loop under some failure scenarios. This value is the initial backoff value and will increase exponentially for each failed request, up to the retry.backoff.max.ms value. | 100 | Integer |
camel.component.kafka-share.sasl-auth-type | Simplified authentication type to use. This provides an easier way to configure Kafka authentication without manually setting securityProtocol, saslMechanism, and saslJaasConfig. When set, the appropriate security settings are automatically derived. Note: This is optional. You can still use the traditional approach with explicit securityProtocol, saslMechanism, and saslJaasConfig properties. | KafkaAuthType | |
camel.component.kafka-share.sasl-jaas-config | Expose the kafka sasl.jaas.config parameter Example: org.apache.kafka.common.security.plain.PlainLoginModule required username=USERNAME password=PASSWORD; | String | |
camel.component.kafka-share.sasl-kerberos-service-name | The Kerberos principal name that Kafka runs as. This can be defined either in Kafka’s JAAS config or in Kafka’s config. | String | |
camel.component.kafka-share.sasl-mechanism | The Simple Authentication and Security Layer (SASL) Mechanism used. For the valid values see http://www.iana.org/assignments/sasl-mechanisms/sasl-mechanisms.xhtml | GSSAPI | String |
camel.component.kafka-share.sasl-password | Password for SASL authentication. Used when saslAuthType is set to PLAIN, SCRAM_SHA_256, or SCRAM_SHA_512. | String | |
camel.component.kafka-share.sasl-username | Username for SASL authentication. Used when saslAuthType is set to PLAIN, SCRAM_SHA_256, or SCRAM_SHA_512. | String | |
camel.component.kafka-share.schema-registry-u-r-l | URL of the schema registry servers to use. The format is host1:port1,host2:port2. This is known as schema.registry.url in multiple Schema registries documentation. This option is only available externally (not standard Apache Kafka) | String | |
camel.component.kafka-share.security-protocol | Protocol used to communicate with brokers. SASL_PLAINTEXT, PLAINTEXT, SASL_SSL and SSL are supported | PLAINTEXT | String |
camel.component.kafka-share.send-buffer-bytes | Socket write buffer size | 131072 | Integer |
camel.component.kafka-share.shutdown-timeout | Timeout in milliseconds to wait gracefully for the consumer or producer to shut down and terminate its worker threads. | 30000 | Integer |
camel.component.kafka-share.specific-avro-reader | This enables the use of a specific Avro reader for use with the in multiple Schema registries documentation with Avro Deserializers implementation. This option is only available externally (not standard Apache Kafka) | false | Boolean |
camel.component.kafka-share.ssl-cipher-suites | A list of cipher suites. This is a named combination of authentication, encryption, MAC and key exchange algorithm used to negotiate the security settings for a network connection using TLS or SSL network protocol. By default, all the available cipher suites are supported. | String | |
camel.component.kafka-share.ssl-context-parameters | SSL configuration using a Camel SSLContextParameters object. If configured, it’s applied before the other SSL endpoint parameters. NOTE: Kafka only supports loading keystore from file locations, so prefix the location with file: in the KeyStoreParameters.resource option. The option is a org.apache.camel.support.jsse.SSLContextParameters type. | SSLContextParameters | |
camel.component.kafka-share.ssl-enabled-protocols | The list of protocols enabled for SSL connections. The default is TLSv1.2,TLSv1.3 when running with Java 11 or newer, TLSv1.2 otherwise. With the default value for Java 11, clients and servers will prefer TLSv1.3 if both support it and fallback to TLSv1.2 otherwise (assuming both support at least TLSv1.2). This default should be fine for most cases. Also see the config documentation for SslProtocol. | String | |
camel.component.kafka-share.ssl-endpoint-algorithm | The endpoint identification algorithm to validate server hostname using server certificate. Use none or false to disable server hostname verification. | https | String |
camel.component.kafka-share.ssl-key-password | The password of the private key in the key store file or the PEM key specified in sslKeystoreKey. This is required for clients only if two-way authentication is configured. | String | |
camel.component.kafka-share.ssl-keymanager-algorithm | The algorithm used by key manager factory for SSL connections. Default value is the key manager factory algorithm configured for the Java Virtual Machine. | SunX509 | String |
camel.component.kafka-share.ssl-keystore-location | The location of the key store file. This is optional for the client and can be used for two-way authentication for the client. | String | |
camel.component.kafka-share.ssl-keystore-password | The store password for the key store file. This is optional for the client and only needed if sslKeystoreLocation is configured. Key store password is not supported for PEM format. | String | |
camel.component.kafka-share.ssl-keystore-type | The file format of the key store file. This is optional for the client. The default value is JKS | JKS | String |
camel.component.kafka-share.ssl-protocol | The SSL protocol used to generate the SSLContext. The default is TLSv1.3 when running with Java 11 or newer, TLSv1.2 otherwise. This value should be fine for most use cases. Allowed values in recent JVMs are TLSv1.2 and TLSv1.3. TLS, TLSv1.1, SSL, SSLv2 and SSLv3 may be supported in older JVMs, but their usage is discouraged due to known security vulnerabilities. With the default value for this config and sslEnabledProtocols, clients will downgrade to TLSv1.2 if the server does not support TLSv1.3. If this config is set to TLSv1.2, clients will not use TLSv1.3 even if it is one of the values in sslEnabledProtocols and the server only supports TLSv1.3. | String | |
camel.component.kafka-share.ssl-provider | The name of the security provider used for SSL connections. Default value is the default security provider of the JVM. | String | |
camel.component.kafka-share.ssl-trustmanager-algorithm | The algorithm used by trust manager factory for SSL connections. Default value is the trust manager factory algorithm configured for the Java Virtual Machine. | PKIX | String |
camel.component.kafka-share.ssl-truststore-location | The location of the trust store file. | String | |
camel.component.kafka-share.ssl-truststore-password | The password for the trust store file. If a password is not set, trust store file configured will still be used, but integrity checking is disabled. Trust store password is not supported for PEM format. | String | |
camel.component.kafka-share.ssl-truststore-type | The file format of the trust store file. The default value is JKS. | JKS | String |
camel.component.kafka-share.use-global-ssl-context-parameters | Enable usage of global SSL context parameters. | false | Boolean |
camel.component.kafka-share.value-deserializer | Deserializer class for value that implements the Deserializer interface. | org.apache.kafka.common.serialization.StringDeserializer | String |