
Apache Kafka TopicCreatePartitions
CertifiedIncrease the partition count of a Kafka topic
Apache Kafka TopicCreatePartitions
Increase the partition count of a Kafka topic
Grows a topic to a new total partition count using the Kafka AdminClient. Partition counts can only be increased, never decreased.
Fails with UnknownTopicOrPartitionException if the topic does not exist, or InvalidPartitionsException if totalPartitionCount is not greater than the current count.
type: io.kestra.plugin.kafka.TopicCreatePartitionsExamples
Scale out a tenant topic to 12 partitions
id: kafka_topic_create_partitions
namespace: company.team
tasks:
- id: create_partitions
type: io.kestra.plugin.kafka.TopicCreatePartitions
properties:
bootstrap.servers: localhost:9092
topic: tenant_acme_orders
totalPartitionCount: 12
Properties
properties *object
Kafka AdminClient properties
Must include bootstrap.servers; accepts any Kafka AdminClient config. Provide base64-encoded content for ssl.keystore.location and ssl.truststore.location when using SSL.
topic *string
Topic name
totalPartitionCount *integerstring
New total partition count
Must be greater than the topic's current partition count.
assets
Assets this task consumes as inputs or produces as outputs, for lineage tracking and the asset graph (Enterprise Edition). A flow declaring this property on a task is rejected in the open-source edition.
io.kestra.core.models.assets.AssetsDeclaration
IGNOREFAILWARNAsset failure behavior
Behavior applied to the task state when a declared asset fails to render, emit, or be persisted (e.g. a lock conflict): FAIL escalates it to FAILED, WARN (default) warns it if it would otherwise succeed, IGNORE leaves the state untouched.
Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.
The assets consumed as inputs.
io.kestra.core.models.assets.AssetIdentifier
1The assets produced as outputs.
io.kestra.plugin.ee.assets.Dataset
1150{}1150io.kestra.plugin.ee.assets.File
1150{}1150io.kestra.plugin.ee.assets.Table
1150{}1150io.kestra.plugin.ee.assets.VM
1150{}1150io.kestra.core.models.assets.External
1150{}1150io.kestra.core.models.assets.Custom
11501Custom asset type
{}1150callTimeout string
PT30SAdminClient call timeout
Maximum duration to wait for each AdminClient operation to complete before failing the task. Defaults to PT30S (30 seconds). Distinct from the task-level timeout, which lets the worker kill the task without retrying it.
Outputs
topic string
Topic name
totalPartitionCount integer
New total partition count