Apache Kafka TopicUpdate

Apache Kafka TopicUpdate

Certified

Update a Kafka topic's configuration

Incrementally alters an existing topic's configuration, typically retention.ms or retention.bytes to enforce a per-tenant data retention policy. Only the provided configs are changed; other configs are left untouched. Fails with UnknownTopicOrPartitionException if the topic does not exist.

yaml
type: io.kestra.plugin.kafka.TopicUpdate

Shrink retention for a tenant topic to 3 days

yaml
id: kafka_topic_update
namespace: company.team

tasks:
  - id: update_topic
    type: io.kestra.plugin.kafka.TopicUpdate
    properties:
      bootstrap.servers: localhost:9092
    topic: tenant_acme_orders
    retentionMs: 259200000
Properties

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 name

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.

Definitions
assetFailureBehaviorstring
Possible Values
IGNOREFAILWARN

Asset 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.

enableAutobooleanstring

Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.

inputsarray

The assets consumed as inputs.

id*string
Min length1
typestring
outputs

The assets produced as outputs.

id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*string
Min length1

Custom asset type

descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
DefaultPT30S

AdminClient 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.

Default{}

Additional topic-level configs to set

Any other Kafka topic config, for example cleanup.policy or min.insync.replicas.

Retention size in bytes

Maps to the topic-level retention.bytes config.

Retention duration in milliseconds

Maps to the topic-level retention.ms config.

Updated topic name

SubTypestring

Configs that were set on the topic