tKafkaInput properties for Apache Spark Structured Streaming | Talend Components for Jobs Help
Skip to main content Skip to complementary content

tKafkaInput properties for Apache Spark Structured Streaming

Last updated: 9/30/2026

These properties are used to configure tKafkaInput running in the Spark Structured Streaming Job framework.

The Spark Structured Streaming tKafkaInput component belongs to the Messaging family.

The streaming version of this component is available in Talend Real-Time Big Data Platform and in Talend Data Fabric.

Basic settings

Properties Description
Schema and Edit schema

A schema is a row description. It defines the number of fields (columns) to be processed and passed on to the next component. When you create a Spark Job, avoid the reserved word line when naming the fields.

  • Built-In: You create and store the schema locally for this component only.

  • Repository: You have already created the schema and stored it in the Repository. You can reuse it in various projects and Job designs.

Click Edit schema to make changes to the schema. If you make changes, the schema automatically becomes built-in.

  • View schema: choose this option to view the schema only.

  • Change to built-in property: choose this option to change the schema to Built-in for local changes.

  • Update repository connection: choose this option to change the schema stored in the repository and decide whether to propagate the changes to all the Jobs upon completion.

    If you just want to propagate the changes to the current Job, you can select No upon completion and choose this schema metadata again in the Repository Content window.

Configuration component Select the tKafkaConfiguration component to use for the Kafka connection settings.
Starting offset

Select the starting point from which the messages of a topic are consumed.

In Kafka, the sequential ID number of a message is called offset. From this list, you can select From beginning to start consumption from the oldest message of the entire topic, or select From latest to start from the latest message that has been consumed by the same consumer group and of which the offset is tracked by Spark within Spark checkpoints.

Each consumer group has its own counter to remember the position of a message it has consumed. For this reason, once a consumer group starts to consume messages of a given topic, a consumer group recognizes the latest message only with regard to the position where this group stops the consumption, rather than to the entire topic. Based on this principle, the following behaviors can be expected:

  • A topic has for example 100 messages. If a consumer group has stopped the consumption at the message of the offset 50, then when you select From latest, the same consumer group restarts from the offset 51.

  • If you create a new consumer group or reset an existing consumer group, which, in either case, means this group has not consumed any message of this topic, then when you start it from latest, this new group starts and waits for the offset 101.

Topic name

Enter the name of the topic from which tKafkaInput receives the feed of messages.

Group ID

Enter the name of the consumer group to which you want the current consumer (the tKafkaInput component) to belong.

This consumer group will be created at runtime if it does not exist at that moment.

This property is available only when you are using Spark 2.0 or the Hadoop distribution to be used is running Spark 2.0. If you do not know the Spark version you are using, ask the administrator of your cluster for details.

Set number of records per second to read from each Kafka partition Select this check box to set the maximum number of records read from each Kafka partition per trigger. In the field that appears, enter the maximum number of records.
Enable watermarking Select this check box to activate watermarking and select the time mode to be used:
  • Event time: Select to use a timestamp column representing when each event occurred, then select the column in the Watermark column (timestamp) drop-down list.
  • Processing time: Select to use the time when Spark reads each record.

In the Watermark delay field, specify how long Spark waits for late data before finalizing each window.

Advanced settings

Properties Description
Do not fail on data loss Select this check box to allow the streaming query to continue processing when Kafka offsets are missing — for example, because topics were deleted or offsets are out of range. If cleared, the query stops with an error.
Kafka properties

Add the Kafka consumer properties you need to customize to this table. For example, you can set a specific zookeeper.connection.timeout.ms value to avoid ZkTimeoutException.

For further information about the consumer properties you can define in this table, see the section describing the consumer configuration in Kafka's documentation in http://kafka.apache.org/documentation.html#consumerconfigs.

Use hierarchical mode

Select this check box to map the binary (including hierarchical) Avro schema to the flat schema defined in the schema editor of the current component. If the Avro message to be processed is flat, leave this check box clear.

Once selecting it, you need set the following parameter(s):

  • Local path to the avro schema: browse to the file which defines the schema of the Avro data to be processed.

  • Mapping: create the map between the schema columns of the current component and the data stored in the hierarchical Avro message to be handled. In the Node column, you need to enter the JSON path pointing to the data to be read from the Avro message.

Usage

Usage guidance Description
Usage rule

This component is used as a start component and requires an output link.

This component, along with the Spark Structured Streaming component Palette it belongs to, appears only when you are creating a Spark Structured Streaming Job.

Spark Connection

You need to use the Structured Streaming Configuration tab in the Run view to define the connection to a Spark cluster for the whole Job.

This connection is effective on a per-Job basis.

Did this page help you?

If you find any issues with this page or its content – a typo, a missing step, or a technical error – please let us know!