This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new db91e00869 Issue #8508 : Add a maximum consume duration to the Kafka
Consumer (#8518)
db91e00869 is described below
commit db91e00869270ef53330a988fa81fa6ee48c9d27
Author: Matt Casters <[email protected]>
AuthorDate: Tue Sep 22 10:02:49 2026 +0200
Issue #8508 : Add a maximum consume duration to the Kafka Consumer (#8518)
---
.../ROOT/pages/how-to-guides/cdc-log-sniffing.adoc | 1 +
.../pages/pipeline/transforms/kafkaconsumer.adoc | 9 +-
.../0005-kafka-consumer-read-max-duration.hpl | 122 +++++++++++++++
.../kafka/main-0005-kafka-test-max-duration.hwf | 171 +++++++++++++++++++++
.../kafka/prepare-kafka-test-0005-max-duration.hpl | 114 ++++++++++++++
.../kafka/consumer/KafkaConsumerInput.java | 160 +++++++++++++++++--
.../kafka/consumer/KafkaConsumerInputData.java | 6 +-
.../kafka/consumer/KafkaConsumerInputDialog.java | 27 +++-
.../kafka/consumer/KafkaConsumerInputMeta.java | 29 ++++
.../consumer/messages/messages_en_US.properties | 4 +
.../kafka/consumer/KafkaConsumerInputMetaTest.java | 106 +++++++++++++
.../kafka/consumer/KafkaConsumerInputTest.java | 60 ++++++++
12 files changed, 793 insertions(+), 16 deletions(-)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/cdc-log-sniffing.adoc
b/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/cdc-log-sniffing.adoc
index a1e10dde54..eaefc6f9a4 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/cdc-log-sniffing.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/how-to-guides/cdc-log-sniffing.adoc
@@ -131,6 +131,7 @@ The xref:pipeline/transforms/kafkaconsumer.adoc[Kafka
consumer] runs a sub-pipel
* *Long-lived*: leave the consumer running and process batches by *Duration*
and/or *Number of records*. Use this for near-real-time apply.
* *Scheduled drain*: enable *Stop when idle* (and an idle timeout) so a
workflow can start the pipeline, empty the topic, and finish. That is the
“queue + schedule” pattern.
+* *Timed window*: set *Max consume duration (ms)* when the topic never goes
idle (a low continuous trickle) but you still want a scheduled run to process
for a fixed time and then finish.
Prefer *Offset management: when batch completed* so a failed sub-pipeline does
not commit past records it never applied.
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
index 70e220097b..33baddc302 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
@@ -40,9 +40,10 @@ A consumer group is a set of consumers sharing a common
group identifier.
By default the Kafka consumer transform continuously ingests streaming data.
To stop after the topic has been drained, enable *Stop when idle* on the Batch
tab (with an optional max idle time).
+To consume for a fixed wall-clock window even when messages keep arriving, set
*Max consume duration (ms)* (for example 15 minutes once per hour).
You can also use the Abort transform in your parent or sub-pipeline to stop
consuming records for other conditions.
-For example, you can run the parent pipeline on a timed schedule, drain a
topic once per hour with stop-when-idle, or abort the sub-pipeline if sensor
data exceeds a preset range.
+For example, you can run the parent pipeline on a timed schedule, drain a
topic once per hour with stop-when-idle, consume for a fixed duration on a
topic that never goes idle, or abort the sub-pipeline if sensor data exceeds a
preset range.
== Options
@@ -94,6 +95,12 @@ While this option is enabled, the consumer polls with a
short (100 ms) timeout s
|Max idle time (ms)|The maximum time in milliseconds to wait without receiving
records before stopping when *Stop when idle* is enabled.
Defaults to 500.
Supports variables.
+|Max consume duration (ms)|Maximum time in milliseconds to consume since the
transform started.
+This is not the batch *Duration (ms)* field: batch duration only controls how
long each poll waits for a batch.
+When this limit is reached the transform finishes gracefully after the current
batch is processed and committed.
+0 or empty means no limit (the default).
+*Stop when idle* can still finish earlier if the topic goes quiet first.
+Supports variables.
|Offset management a|Choose when to commit
* when record read
diff --git a/integration-tests/kafka/0005-kafka-consumer-read-max-duration.hpl
b/integration-tests/kafka/0005-kafka-consumer-read-max-duration.hpl
new file mode 100644
index 0000000000..dc692d9455
--- /dev/null
+++ b/integration-tests/kafka/0005-kafka-consumer-read-max-duration.hpl
@@ -0,0 +1,122 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+Licensed to the Apache Software Foundation (ASF) under one or more
+contributor license agreements. See the NOTICE file distributed with
+this work for additional information regarding copyright ownership.
+The ASF licenses this file to You under the Apache License, Version 2.0
+(the "License"); you may not use this file except in compliance with
+the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+
+-->
+<pipeline>
+ <info>
+ <name>0005-kafka-consumer-read-max-duration</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <pipeline_status>0</pipeline_status>
+ <parameters>
+ </parameters>
+ <capture_transform_performance>N</capture_transform_performance>
+
<transform_performance_capturing_delay>1000</transform_performance_capturing_delay>
+
<transform_performance_capturing_size_limit>100</transform_performance_capturing_size_limit>
+ <created_user>-</created_user>
+ <created_date>2021/12/21 09:42:58.199</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2021/12/21 09:42:58.199</modified_date>
+ <key_for_session_key>H4sIAAAAAAAAAAMAAAAAAAAAAAA=</key_for_session_key>
+ <is_key_private>N</is_key_private>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>Kafka Consumer</from>
+ <to>Log Output</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>Kafka Consumer</name>
+ <type>KafkaConsumer</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <topic>hop-test-max-duration</topic>
+ <consumerGroup>hop-max-duration</consumerGroup>
+
<pipelinePath>${PROJECT_HOME}/0001-kafka-consumer-called-subpipeline.hpl</pipelinePath>
+ <subTransform>Out record</subTransform>
+ <batchSize>10</batchSize>
+ <batchDuration>500</batchDuration>
+ <stopWhenIdle>N</stopWhenIdle>
+ <maxIdleTimeMs>500</maxIdleTimeMs>
+ <maxConsumeDurationMs>10000</maxConsumeDurationMs>
+ <directBootstrapServers>${BOOTSTRAP_SERVERS}</directBootstrapServers>
+ <AUTO_COMMIT>N</AUTO_COMMIT>
+ <OutputField kafkaName="key" type="String">Key</OutputField>
+ <OutputField kafkaName="message" type="String">Message</OutputField>
+ <OutputField kafkaName="topic" type="String">Topic</OutputField>
+ <OutputField kafkaName="partition" type="Integer">Partition</OutputField>
+ <OutputField kafkaName="offset" type="Integer">Offset</OutputField>
+ <OutputField kafkaName="timestamp" type="Integer">Timestamp</OutputField>
+ <advancedConfig>
+ <option property="auto.offset.reset" value="earliest"/>
+ <option property="ssl.key.password" value=""/>
+ <option property="ssl.keystore.location" value=""/>
+ <option property="ssl.keystore.password" value=""/>
+ <option property="ssl.truststore.location" value=""/>
+ <option property="ssl.truststore.password" value=""/>
+ </advancedConfig>
+ <attributes/>
+ <GUI>
+ <xloc>448</xloc>
+ <yloc>112</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Log Output</name>
+ <type>WriteToLog</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <loglevel>log_level_basic</loglevel>
+ <displayHeader>Y</displayHeader>
+ <limitRows>N</limitRows>
+ <limitRowsNumber>0</limitRowsNumber>
+ <logmessage>Kafka Output</logmessage>
+ <fields>
+ <field>
+ <name>message</name>
+ </field>
+ </fields>
+ <attributes/>
+ <GUI>
+ <xloc>560</xloc>
+ <yloc>112</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git a/integration-tests/kafka/main-0005-kafka-test-max-duration.hwf
b/integration-tests/kafka/main-0005-kafka-test-max-duration.hwf
new file mode 100644
index 0000000000..ce3f4c30ff
--- /dev/null
+++ b/integration-tests/kafka/main-0005-kafka-test-max-duration.hwf
@@ -0,0 +1,171 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+Licensed to the Apache Software Foundation (ASF) under one or more
+contributor license agreements. See the NOTICE file distributed with
+this work for additional information regarding copyright ownership.
+The ASF licenses this file to You under the Apache License, Version 2.0
+(the "License"); you may not use this file except in compliance with
+the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+
+-->
+<workflow>
+ <name>main-0005-kafka-test-max-duration</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <created_user>-</created_user>
+ <modified_user>-</modified_user>
+ <created_date>2021/12/30 10:42:13.629</created_date>
+ <modified_date>2021/12/30 10:42:13.629</modified_date>
+ <workflow_version/>
+ <parameters>
+ <parameter>
+ <name>BOOTSTRAP_SERVERS</name>
+ <description/>
+ <default_value>kafka:9092</default_value>
+ </parameter>
+ </parameters>
+ <actions>
+ <action>
+ <repeat>N</repeat>
+ <schedulerType>0</schedulerType>
+ <intervalSeconds>0</intervalSeconds>
+ <intervalMinutes>60</intervalMinutes>
+ <DayOfMonth>1</DayOfMonth>
+ <weekDay>1</weekDay>
+ <minutes>0</minutes>
+ <hour>12</hour>
+ <doNotWaitOnFirstExecution>N</doNotWaitOnFirstExecution>
+ <name>Start</name>
+ <description/>
+ <type>SPECIAL</type>
+ <attributes/>
+ <xloc>128</xloc>
+ <yloc>144</yloc>
+ <parallel>N</parallel>
+ <attributes_hac/>
+ </action>
+ <action>
+
<filename>${PROJECT_HOME}/prepare-kafka-test-0005-max-duration.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <create_parent_folder>N</create_parent_folder>
+ <set_logfile>N</set_logfile>
+ <set_append_logfile>N</set_append_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <wait_until_finished>Y</wait_until_finished>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <run_configuration>local</run_configuration>
+ <name>prepare-kafka-test-0005-max-duration.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <xloc>336</xloc>
+ <yloc>144</yloc>
+ <parallel>N</parallel>
+ <attributes_hac/>
+ </action>
+ <action>
+ <script>
+var txt = previous_result.getLogText();
+
+var ok = false;
+
+var expectedValues = [
+ "Message = Hello Hop Duration!",
+ ];
+
+for (var i = 0 ; i<expectedValues.length ; i++) {
+ var expectedValue = expectedValues[i];
+ if (txt.contains(expectedValue)) {
+ ok = true;
+ log.logBasic("Expected value logged as ''" + expectedValue + "'' FOUND!");
+ }
+}
+
+if (!ok) {
+ log.logBasic("Expected value logged as ''" + expectedValues[0] + "'' NOT
FOUND!");
+}
+
+ok;</script>
+ <name>Check log</name>
+ <description/>
+ <type>EVAL</type>
+ <attributes/>
+ <xloc>864</xloc>
+ <yloc>144</yloc>
+ <parallel>N</parallel>
+ <attributes_hac/>
+ </action>
+ <action>
+
<filename>${PROJECT_HOME}/0005-kafka-consumer-read-max-duration.hpl</filename>
+ <params_from_previous>N</params_from_previous>
+ <exec_per_row>N</exec_per_row>
+ <clear_rows>N</clear_rows>
+ <clear_files>N</clear_files>
+ <create_parent_folder>N</create_parent_folder>
+ <set_logfile>N</set_logfile>
+ <set_append_logfile>N</set_append_logfile>
+ <logfile/>
+ <logext/>
+ <add_date>N</add_date>
+ <add_time>N</add_time>
+ <loglevel>Basic</loglevel>
+ <wait_until_finished>Y</wait_until_finished>
+ <parameters>
+ <pass_all_parameters>Y</pass_all_parameters>
+ </parameters>
+ <run_configuration>local</run_configuration>
+ <name>0005-kafka-consumer-read-max-duration.hpl</name>
+ <description/>
+ <type>PIPELINE</type>
+ <attributes/>
+ <xloc>608</xloc>
+ <yloc>144</yloc>
+ <parallel>N</parallel>
+ <attributes_hac/>
+ </action>
+ </actions>
+ <hops>
+ <hop>
+ <from>Start</from>
+ <to>prepare-kafka-test-0005-max-duration.hpl</to>
+ <evaluation>Y</evaluation>
+ <unconditional>Y</unconditional>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>0005-kafka-consumer-read-max-duration.hpl</from>
+ <to>Check log</to>
+ <evaluation>N</evaluation>
+ <unconditional>Y</unconditional>
+ <enabled>Y</enabled>
+ </hop>
+ <hop>
+ <from>prepare-kafka-test-0005-max-duration.hpl</from>
+ <to>0005-kafka-consumer-read-max-duration.hpl</to>
+ <evaluation>N</evaluation>
+ <unconditional>Y</unconditional>
+ <enabled>Y</enabled>
+ </hop>
+ </hops>
+ <notepads/>
+ <attributes/>
+</workflow>
diff --git a/integration-tests/kafka/prepare-kafka-test-0005-max-duration.hpl
b/integration-tests/kafka/prepare-kafka-test-0005-max-duration.hpl
new file mode 100644
index 0000000000..26db725d97
--- /dev/null
+++ b/integration-tests/kafka/prepare-kafka-test-0005-max-duration.hpl
@@ -0,0 +1,114 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+Licensed to the Apache Software Foundation (ASF) under one or more
+contributor license agreements. See the NOTICE file distributed with
+this work for additional information regarding copyright ownership.
+The ASF licenses this file to You under the Apache License, Version 2.0
+(the "License"); you may not use this file except in compliance with
+the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+
+-->
+<pipeline>
+ <info>
+ <name>prepare-kafka-test-0005-max-duration</name>
+ <name_sync_with_filename>Y</name_sync_with_filename>
+ <description/>
+ <extended_description/>
+ <pipeline_version/>
+ <pipeline_type>Normal</pipeline_type>
+ <parameters>
+ </parameters>
+ <capture_transform_performance>N</capture_transform_performance>
+
<transform_performance_capturing_delay>1000</transform_performance_capturing_delay>
+
<transform_performance_capturing_size_limit>100</transform_performance_capturing_size_limit>
+ <created_user>-</created_user>
+ <created_date>2021/12/21 09:37:38.673</created_date>
+ <modified_user>-</modified_user>
+ <modified_date>2021/12/21 09:37:38.673</modified_date>
+ <key_for_session_key>H4sIAAAAAAAAAAMAAAAAAAAAAAA=</key_for_session_key>
+ <is_key_private>N</is_key_private>
+ </info>
+ <notepads>
+ </notepads>
+ <order>
+ <hop>
+ <from>Generate rows</from>
+ <to>Kafka Producer</to>
+ <enabled>Y</enabled>
+ </hop>
+ </order>
+ <transform>
+ <name>Generate rows</name>
+ <type>RowGenerator</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <fields>
+ <field>
+ <length>-1</length>
+ <name>field</name>
+ <precision>-1</precision>
+ <set_empty_string>N</set_empty_string>
+ <type>String</type>
+ <nullif>Hello Hop Duration!</nullif>
+ </field>
+ </fields>
+ <interval_in_ms>5000</interval_in_ms>
+ <last_time_field>FiveSecondsAgo</last_time_field>
+ <never_ending>N</never_ending>
+ <limit>3</limit>
+ <row_time_field>now</row_time_field>
+ <attributes/>
+ <GUI>
+ <xloc>288</xloc>
+ <yloc>160</yloc>
+ </GUI>
+ </transform>
+ <transform>
+ <name>Kafka Producer</name>
+ <type>KafkaProducerOutput</type>
+ <description/>
+ <distribute>Y</distribute>
+ <custom_distribution/>
+ <copies>1</copies>
+ <partitioning>
+ <method>none</method>
+ <schema_name/>
+ </partitioning>
+ <directBootstrapServers>${BOOTSTRAP_SERVERS}</directBootstrapServers>
+ <topic>hop-test-max-duration</topic>
+ <clientId/>
+ <keyField/>
+ <messageField>field</messageField>
+ <advancedConfig>
+ <option property="compression.type" value="none"/>
+ <option property="ssl.key.password" value=""/>
+ <option property="ssl.keystore.location" value=""/>
+ <option property="ssl.keystore.password" value=""/>
+ <option property="ssl.truststore.location" value=""/>
+ <option property="ssl.truststore.password" value=""/>
+ </advancedConfig>
+ <attributes/>
+ <GUI>
+ <xloc>464</xloc>
+ <yloc>160</yloc>
+ </GUI>
+ </transform>
+ <transform_error_handling>
+ </transform_error_handling>
+ <attributes/>
+</pipeline>
diff --git
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
index 59970c7610..75945faa39 100644
---
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
+++
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
@@ -88,7 +88,19 @@ public class KafkaConsumerInput
data.batchSize = Const.toIntExpanded(resolve(meta.getBatchSize()), 0);
data.stopWhenIdle = meta.isStopWhenIdle();
data.maxIdleTimeMs = Const.toLong(resolve(meta.getMaxIdleTimeMs()), 500L);
- data.lastRecordTime = System.currentTimeMillis();
+ long maxConsume = Const.toLong(resolve(meta.getMaxConsumeDurationMs()),
0L);
+ data.maxConsumeDurationMs = maxConsume > 0 ? maxConsume : 0L;
+ data.startTime = System.currentTimeMillis();
+ data.lastRecordTime = data.startTime;
+ logBasic(
+ "Kafka consumer batchDuration="
+ + data.batchDuration
+ + "ms, stopWhenIdle="
+ + data.stopWhenIdle
+ + ", maxIdleTimeMs="
+ + data.maxIdleTimeMs
+ + ", maxConsumeDurationMs="
+ + data.maxConsumeDurationMs);
data.consumer = buildKafkaConsumer(this, meta);
@@ -108,6 +120,7 @@ public class KafkaConsumerInput
// Set Kafka consumer is closing flag to false
data.isKafkaConsumerClosing = false;
+ startMaxConsumeDeadlineWakeup();
return true;
}
@@ -212,7 +225,9 @@ public class KafkaConsumerInput
@Override
public void dispose() {
+ interruptMaxConsumeDeadlineWakeup();
if (data.consumer != null) {
+ data.consumer.wakeup();
data.consumer.unsubscribe();
data.consumer.close();
}
@@ -284,16 +299,39 @@ public class KafkaConsumerInput
// Poll records...
// If we get any, process them...
// When stop-when-idle is enabled, use a short poll timeout so idle time
can be measured.
+ // When a max consume duration is set, cap the poll to the remaining time
so a long batch
+ // duration cannot overshoot the deadline.
//
try {
+ long now = System.currentTimeMillis();
+ if (maxConsumeDurationReached(now, data.maxConsumeDurationMs,
data.startTime)) {
+ return stopGracefully(
+ "Kafka consumer max consume duration of "
+ + data.maxConsumeDurationMs
+ + "ms reached, stopping gracefully");
+ }
long pollMs =
- data.stopWhenIdle ? 100L : (data.batchDuration > 0 ?
data.batchDuration : Long.MAX_VALUE);
+ pollTimeoutMs(
+ data.stopWhenIdle,
+ data.batchDuration,
+ data.maxConsumeDurationMs,
+ data.startTime,
+ now);
Duration duration = Duration.ofMillis(pollMs);
ConsumerRecords<Object, Object> records = data.consumer.poll(duration);
if (!data.isKafkaConsumerClosing) {
if (records.isEmpty()) {
- // No records: optionally stop after max idle time.
+ // No records: still honor max consume duration. The deadline is
wall-clock since
+ // start, not "time since last message", so an idle topic must stop
here.
+ if (maxConsumeDurationReached(
+ System.currentTimeMillis(), data.maxConsumeDurationMs,
data.startTime)) {
+ return stopGracefully(
+ "Kafka consumer max consume duration of "
+ + data.maxConsumeDurationMs
+ + "ms reached, stopping gracefully");
+ }
+ // Optionally stop after max idle time.
// Do not count idle until partitions are assigned — group join /
rebalance can take
// longer than maxIdleTimeMs and would otherwise stop before any
poll can succeed.
//
@@ -301,16 +339,10 @@ public class KafkaConsumerInput
if (data.consumer.assignment() == null ||
data.consumer.assignment().isEmpty()) {
data.lastRecordTime = System.currentTimeMillis();
} else if ((System.currentTimeMillis() - data.lastRecordTime) >=
data.maxIdleTimeMs) {
- logBasic(
+ return stopGracefully(
"Kafka consumer idle timeout of "
+ data.maxIdleTimeMs
+ "ms exceeded, stopping gracefully");
- data.isKafkaConsumerClosing = true;
- if (data.executor != null) {
- data.executor.getPipeline().stopAll();
- }
- setOutputDone();
- return false;
}
}
} else {
@@ -372,13 +404,31 @@ public class KafkaConsumerInput
data.incomingRowsBuffer.clear();
}
}
+
+ if (maxConsumeDurationReached(
+ System.currentTimeMillis(), data.maxConsumeDurationMs,
data.startTime)) {
+ return stopGracefully(
+ "Kafka consumer max consume duration of "
+ + data.maxConsumeDurationMs
+ + "ms reached, stopping gracefully");
+ }
}
} catch (WakeupException e) {
- // We're going to close kafka consumer because of pipeline has been
stopped so stop executor
- // too
- data.executor.getPipeline().stopAll();
+ // Deadline wakeup (no new messages, poll was still blocked) or the
pipeline was stopped.
+ if (data.maxConsumeDeadlineWakeup
+ || maxConsumeDurationReached(
+ System.currentTimeMillis(), data.maxConsumeDurationMs,
data.startTime)) {
+ return stopGracefully(
+ "Kafka consumer max consume duration of "
+ + data.maxConsumeDurationMs
+ + "ms reached, stopping gracefully");
+ }
+ if (data.executor != null) {
+ data.executor.getPipeline().stopAll();
+ }
setOutputDone();
stopAll();
+ return false;
}
if (data.executor.getErrors() > 0 && errorHandlingConditionIsSatisfied()) {
@@ -401,6 +451,90 @@ public class KafkaConsumerInput
return true;
}
+ /**
+ * True when a max consume duration is configured and the wall clock since
transform start has
+ * reached it. {@code maxConsumeDurationMs <= 0} means no limit.
+ */
+ static boolean maxConsumeDurationReached(long now, long
maxConsumeDurationMs, long startTime) {
+ return maxConsumeDurationMs > 0 && (now - startTime) >=
maxConsumeDurationMs;
+ }
+
+ /**
+ * Poll timeout in milliseconds. Stop-when-idle and max-consume-duration use
a short poll so the
+ * deadline can be re-checked when no records arrive. A long or infinite
poll would otherwise
+ * never return on an idle topic, and the duration check after poll() would
never run. A
+ * configured max consume duration also caps the timeout to the remaining
window.
+ */
+ static long pollTimeoutMs(
+ boolean stopWhenIdle,
+ long batchDuration,
+ long maxConsumeDurationMs,
+ long startTime,
+ long now) {
+ boolean shortPoll = stopWhenIdle || maxConsumeDurationMs > 0;
+ long pollMs = shortPoll ? 100L : (batchDuration > 0 ? batchDuration :
Long.MAX_VALUE);
+ if (maxConsumeDurationMs > 0) {
+ long remaining = maxConsumeDurationMs - (now - startTime);
+ if (remaining <= 0) {
+ return 0L;
+ }
+ pollMs = Math.min(pollMs, remaining);
+ }
+ return pollMs;
+ }
+
+ /**
+ * {@code consumer.poll()} can block beyond the requested timeout
(coordinator lookup, metadata, a
+ * stuck fetch). Interrupt that wait when the max consume deadline is
reached so an idle topic
+ * still finishes.
+ */
+ private void startMaxConsumeDeadlineWakeup() {
+ if (data.maxConsumeDurationMs <= 0 || data.consumer == null) {
+ return;
+ }
+ final long deadline = data.startTime + data.maxConsumeDurationMs;
+ data.maxConsumeDeadlineThread =
+ new Thread(
+ () -> {
+ try {
+ long sleepMs = deadline - System.currentTimeMillis();
+ if (sleepMs > 0) {
+ Thread.sleep(sleepMs);
+ }
+ if (!data.isKafkaConsumerClosing && data.consumer != null) {
+ data.maxConsumeDeadlineWakeup = true;
+ data.consumer.wakeup();
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ },
+ "KafkaConsumer-maxConsumeDeadline");
+ data.maxConsumeDeadlineThread.setDaemon(true);
+ data.maxConsumeDeadlineThread.start();
+ }
+
+ private void interruptMaxConsumeDeadlineWakeup() {
+ if (data.maxConsumeDeadlineThread != null) {
+ data.maxConsumeDeadlineThread.interrupt();
+ data.maxConsumeDeadlineThread = null;
+ }
+ }
+
+ private boolean stopGracefully(String reason) {
+ logBasic(reason);
+ data.isKafkaConsumerClosing = true;
+ interruptMaxConsumeDeadlineWakeup();
+ if (data.consumer != null) {
+ data.consumer.wakeup();
+ }
+ if (data.executor != null) {
+ data.executor.getPipeline().stopAll();
+ }
+ setOutputDone();
+ return false;
+ }
+
private boolean errorHandlingConditionIsSatisfied() {
// Added a check to be sure that lines collecting for error handling is
limited
// to the case of batchSize = 1.
diff --git
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputData.java
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputData.java
index f96ea2b88b..b3ade646f0 100644
---
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputData.java
+++
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputData.java
@@ -33,10 +33,14 @@ public class KafkaConsumerInputData extends
BaseTransformData implements ITransf
public int batchSize;
public boolean stopWhenIdle;
public long maxIdleTimeMs;
+ public long maxConsumeDurationMs;
+ public long startTime;
public long lastRecordTime;
public RowProducer rowProducer;
public SingleThreadedPipelineExecutor executor;
- public boolean isKafkaConsumerClosing;
+ public volatile boolean isKafkaConsumerClosing;
+ public volatile boolean maxConsumeDeadlineWakeup;
+ public Thread maxConsumeDeadlineThread;
public List<Object[]> incomingRowsBuffer;
/** */
diff --git
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputDialog.java
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputDialog.java
index 0ef5d7b818..fc3a9f29d0 100644
---
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputDialog.java
+++
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputDialog.java
@@ -109,6 +109,8 @@ public class KafkaConsumerInputDialog extends
BaseTransformDialog {
protected Button wStopWhenIdle;
protected Label wlMaxIdleTimeMs;
protected TextVar wMaxIdleTimeMs;
+ protected Label wlMaxConsumeDurationMs;
+ protected TextVar wMaxConsumeDurationMs;
protected CTabFolder wTabFolder;
protected CTabItem wSetupTab;
@@ -304,6 +306,7 @@ public class KafkaConsumerInputDialog extends
BaseTransformDialog {
m.setBatchDuration(wBatchDuration.getText());
m.setStopWhenIdle(wStopWhenIdle.getSelection());
m.setMaxIdleTimeMs(wMaxIdleTimeMs.getText());
+ m.setMaxConsumeDurationMs(wMaxConsumeDurationMs.getText());
m.setSubTransform(wSubTransform.getText());
setTopicsFromTable();
@@ -334,7 +337,7 @@ public class KafkaConsumerInputDialog extends
BaseTransformDialog {
wOffsetGroup.setLayout(flOffsetGroup);
FormData fdOffsetGroup = new FormData();
- fdOffsetGroup.top = new FormAttachment(wMaxIdleTimeMs, 15);
+ fdOffsetGroup.top = new FormAttachment(wMaxConsumeDurationMs, 15);
fdOffsetGroup.left = new FormAttachment(0, 0);
fdOffsetGroup.right = new FormAttachment(100, 0);
wOffsetGroup.setLayoutData(fdOffsetGroup);
@@ -650,6 +653,27 @@ public class KafkaConsumerInputDialog extends
BaseTransformDialog {
fdMaxIdleTimeMs.top = new FormAttachment(wlMaxIdleTimeMs, 0, SWT.CENTER);
wMaxIdleTimeMs.setLayoutData(fdMaxIdleTimeMs);
+ wlMaxConsumeDurationMs = new Label(wBatchComp, SWT.RIGHT);
+ PropsUi.setLook(wlMaxConsumeDurationMs);
+ wlMaxConsumeDurationMs.setText(
+ BaseMessages.getString(PKG,
"KafkaConsumerInputDialog.MaxConsumeDurationMs"));
+ FormData fdlMaxConsumeDurationMs = new FormData();
+ fdlMaxConsumeDurationMs.left = new FormAttachment(0, 0);
+ fdlMaxConsumeDurationMs.top = new FormAttachment(wMaxIdleTimeMs, margin);
+ fdlMaxConsumeDurationMs.right = new FormAttachment(middle, -margin);
+ wlMaxConsumeDurationMs.setLayoutData(fdlMaxConsumeDurationMs);
+
+ wMaxConsumeDurationMs = new TextVar(variables, wBatchComp, SWT.SINGLE |
SWT.LEFT | SWT.BORDER);
+ PropsUi.setLook(wMaxConsumeDurationMs);
+ wMaxConsumeDurationMs.setToolTipText(
+ BaseMessages.getString(PKG,
"KafkaConsumerInputDialog.MaxConsumeDurationMs.Tooltip"));
+ wMaxConsumeDurationMs.addModifyListener(lsMod);
+ FormData fdMaxConsumeDurationMs = new FormData();
+ fdMaxConsumeDurationMs.left = new FormAttachment(wlMaxConsumeDurationMs,
margin);
+ fdMaxConsumeDurationMs.right = new FormAttachment(100, 0);
+ fdMaxConsumeDurationMs.top = new FormAttachment(wlMaxConsumeDurationMs, 0,
SWT.CENTER);
+ wMaxConsumeDurationMs.setLayoutData(fdMaxConsumeDurationMs);
+
wBatchComp.layout();
wBatchTab.setControl(wBatchComp);
}
@@ -865,6 +889,7 @@ public class KafkaConsumerInputDialog extends
BaseTransformDialog {
wBatchDuration.setText(Const.NVL(meta.getBatchDuration(), ""));
wStopWhenIdle.setSelection(meta.isStopWhenIdle());
wMaxIdleTimeMs.setText(Const.NVL(meta.getMaxIdleTimeMs(), "500"));
+ wMaxConsumeDurationMs.setText(Const.NVL(meta.getMaxConsumeDurationMs(),
"0"));
wbAutoCommit.setSelection(meta.isAutoCommit());
wbManualCommit.setSelection(!meta.isAutoCommit());
diff --git
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMeta.java
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMeta.java
index e1671aefbe..3b7737adb8 100644
---
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMeta.java
+++
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMeta.java
@@ -181,6 +181,12 @@ public class KafkaConsumerInputMeta
injectionKeyDescription =
"KafkaConsumerInputMeta.Injection.MAX_IDLE_TIME_MS")
private String maxIdleTimeMs;
+ @HopMetadataProperty(
+ key = "maxConsumeDurationMs",
+ injectionKey = "MAX_CONSUME_DURATION_MS",
+ injectionKeyDescription =
"KafkaConsumerInputMeta.Injection.MAX_CONSUME_DURATION_MS")
+ private String maxConsumeDurationMs;
+
@HopMetadataProperty(
groupKey = "options",
key = "option",
@@ -209,6 +215,7 @@ public class KafkaConsumerInputMeta
batchSize = "1000";
batchDuration = "1000";
maxIdleTimeMs = "500";
+ maxConsumeDurationMs = "0";
subTransform = "";
topics = new ArrayList<>();
options = new ArrayList<>();
@@ -421,6 +428,28 @@ public class KafkaConsumerInputMeta
transformMeta));
}
}
+
+ String maxConsumeResolved =
variables.resolve(Const.NVL(getMaxConsumeDurationMs(), "0"));
+ if (StringUtils.isNotBlank(maxConsumeResolved)) {
+ try {
+ long maxConsume = Long.parseLong(maxConsumeResolved);
+ if (maxConsume < 0) {
+ remarks.add(
+ new CheckResult(
+ ICheckResult.TYPE_RESULT_ERROR,
+ BaseMessages.getString(
+ PKG, "KafkaConsumerInputMeta.CheckResult.Negative", "Max
consume duration"),
+ transformMeta));
+ }
+ } catch (NumberFormatException e) {
+ remarks.add(
+ new CheckResult(
+ ICheckResult.TYPE_RESULT_ERROR,
+ BaseMessages.getString(
+ PKG, "KafkaConsumerInputMeta.CheckResult.NaN", "Max
consume duration"),
+ transformMeta));
+ }
+ }
}
@Override
diff --git
a/plugins/transforms/kafka/src/main/resources/org/apache/hop/pipeline/transforms/kafka/consumer/messages/messages_en_US.properties
b/plugins/transforms/kafka/src/main/resources/org/apache/hop/pipeline/transforms/kafka/consumer/messages/messages_en_US.properties
index 5703e114e7..14fd296b3e 100644
---
a/plugins/transforms/kafka/src/main/resources/org/apache/hop/pipeline/transforms/kafka/consumer/messages/messages_en_US.properties
+++
b/plugins/transforms/kafka/src/main/resources/org/apache/hop/pipeline/transforms/kafka/consumer/messages/messages_en_US.properties
@@ -25,6 +25,8 @@ KafkaConsumerInputDialog.BatchSize=Number of records
KafkaConsumerInputDialog.BatchTab=Batch
KafkaConsumerInputDialog.BootstrapServers=Bootstrap servers
KafkaConsumerInputDialog.MaxIdleTimeMs=Max idle time (ms)
+KafkaConsumerInputDialog.MaxConsumeDurationMs=Max consume duration (ms)
+KafkaConsumerInputDialog.MaxConsumeDurationMs.Tooltip=Maximum time in
milliseconds to consume since the transform started. 0 or empty means no limit.
Stop when idle can still finish earlier. Supports variables such as
'${MAX_CONSUME_MS}'.
KafkaConsumerInputDialog.StopWhenIdle=Stop when idle
KafkaConsumerInputDialog.Column.Name=Output name
KafkaConsumerInputDialog.Column.Ref=Input name
@@ -52,6 +54,7 @@ KafkaConsumerInputDialog.TopicField=Topic
KafkaConsumerInputDialog.Topics=Topics
KafkaConsumerInputDialog.TransformName.Label=Transform name
KafkaConsumerInputMeta.CheckResult.NaN=The "{0}" field is using a non-numeric
value. Please set a numeric value.
+KafkaConsumerInputMeta.CheckResult.Negative=The "{0}" field cannot be
negative. Use 0 for no limit, or a positive duration in milliseconds.
KafkaConsumerInputMeta.CheckResult.NoBatchDefined=The "Number of records" and
"Duration" fields can’t both be set to 0. Please set a value of 1 or higher for
one of the fields.
KafkaConsumerInputMeta.Injection.PIPELINE_PATH=Pipeline path
@@ -66,6 +69,7 @@
KafkaConsumerInputMeta.Injection.DIRECT_BOOTSTRAP_SERVERS=Specify the Bootstrap
KafkaConsumerInputMeta.Injection.BATCH_DURATION=The amount of time to batch
before consuming the messages.
KafkaConsumerInputMeta.Injection.STOP_WHEN_IDLE=Stop the Kafka consumer when
no records are received for the max idle time.
KafkaConsumerInputMeta.Injection.MAX_IDLE_TIME_MS=The maximum idle time in
milliseconds before stopping when stop when idle is enabled.
+KafkaConsumerInputMeta.Injection.MAX_CONSUME_DURATION_MS=The maximum time in
milliseconds to consume messages since the transform started. 0 means no limit.
KafkaConsumerInputMeta.Injection.KEY.OUTPUT_NAME=The name of the output field
for the key.
KafkaConsumerInputMeta.Injection.KEY.TYPE=Specify the data type for the key:
String, Integer, Binary, or Number.
KafkaConsumerInputMeta.Injection.MESSAGE.OUTPUT_NAME=The name of the output
field for the message.
diff --git
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMetaTest.java
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMetaTest.java
index 4d66927e54..3cba4b8ad4 100644
---
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMetaTest.java
+++
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputMetaTest.java
@@ -22,11 +22,15 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.util.ArrayList;
import java.util.List;
import org.apache.commons.lang3.StringUtils;
+import org.apache.hop.core.ICheckResult;
+import org.apache.hop.core.variables.Variables;
import org.apache.hop.core.xml.XmlHandler;
import org.apache.hop.metadata.serializer.memory.MemoryMetadataProvider;
import org.apache.hop.metadata.serializer.xml.XmlMetadataUtil;
+import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.pipeline.transform.TransformSerializationTestUtil;
import org.apache.hop.pipeline.transforms.kafka.shared.KafkaOption;
import org.junit.jupiter.api.Test;
@@ -106,6 +110,7 @@ class KafkaConsumerInputMetaTest {
meta.setAutoCommit(true);
meta.setStopWhenIdle(true);
meta.setMaxIdleTimeMs("1500");
+ meta.setMaxConsumeDurationMs("900000");
meta.getOptions().clear();
meta.getOptions().add(new KafkaOption("auto.offset.reset", "latest"));
meta.getOptions().add(new KafkaOption("ssl.key.password", ""));
@@ -147,6 +152,7 @@ class KafkaConsumerInputMetaTest {
assertEquals("222", copy.getBatchDuration());
assertTrue(copy.isStopWhenIdle());
assertEquals("1500", copy.getMaxIdleTimeMs());
+ assertEquals("900000", copy.getMaxConsumeDurationMs());
assertTrue(copy.getTopics().contains("topic1"));
assertTrue(copy.getTopics().contains("topic2"));
@@ -200,6 +206,7 @@ class KafkaConsumerInputMetaTest {
// Legacy XML without the new fields keeps defaults
assertFalse(meta.isStopWhenIdle());
assertEquals("500", meta.getMaxIdleTimeMs());
+ assertEquals("0", meta.getMaxConsumeDurationMs());
assertEquals(6, meta.getOptions().size());
assertEquals("auto.offset.reset",
meta.getOptions().getFirst().getProperty());
@@ -241,6 +248,7 @@ class KafkaConsumerInputMetaTest {
KafkaConsumerInputMeta meta = new KafkaConsumerInputMeta();
assertFalse(meta.isStopWhenIdle());
assertEquals("500", meta.getMaxIdleTimeMs());
+ assertEquals("0", meta.getMaxConsumeDurationMs());
}
@Test
@@ -248,10 +256,108 @@ class KafkaConsumerInputMetaTest {
KafkaConsumerInputMeta meta = new KafkaConsumerInputMeta();
meta.setStopWhenIdle(true);
meta.setMaxIdleTimeMs("2500");
+ meta.setMaxConsumeDurationMs("15000");
KafkaConsumerInputMeta copy = (KafkaConsumerInputMeta) meta.clone();
assertTrue(copy.isStopWhenIdle());
assertEquals("2500", copy.getMaxIdleTimeMs());
+ assertEquals("15000", copy.getMaxConsumeDurationMs());
+ }
+
+ @Test
+ void testMaxConsumeDurationCheckNan() {
+ KafkaConsumerInputMeta meta = new KafkaConsumerInputMeta();
+ meta.setMaxConsumeDurationMs("not-a-number");
+ List<ICheckResult> remarks = new ArrayList<>();
+ meta.check(
+ remarks,
+ null,
+ new TransformMeta(),
+ null,
+ new String[0],
+ new String[0],
+ null,
+ new Variables(),
+ null);
+ assertTrue(
+ remarks.stream()
+ .anyMatch(
+ r ->
+ r.getType() == ICheckResult.TYPE_RESULT_ERROR
+ && r.getText().contains("Max consume duration")));
+ }
+
+ @Test
+ void testMaxConsumeDurationCheckNegative() {
+ KafkaConsumerInputMeta meta = new KafkaConsumerInputMeta();
+ meta.setMaxConsumeDurationMs("-1");
+ List<ICheckResult> remarks = new ArrayList<>();
+ meta.check(
+ remarks,
+ null,
+ new TransformMeta(),
+ null,
+ new String[0],
+ new String[0],
+ null,
+ new Variables(),
+ null);
+ assertTrue(
+ remarks.stream()
+ .anyMatch(
+ r ->
+ r.getType() == ICheckResult.TYPE_RESULT_ERROR
+ && r.getText().contains("cannot be negative")));
+ }
+
+ @Test
+ void testMaxConsumeDurationLoadedFromLegacyStyleXml() throws Exception {
+ String xml =
+ """
+ <transform>
+ <topic>hop-test-max-duration</topic>
+ <consumerGroup>hop-max-duration</consumerGroup>
+ <pipelinePath>${PROJECT_HOME}/child.hpl</pipelinePath>
+ <subTransform>Out record</subTransform>
+ <batchSize>10</batchSize>
+ <batchDuration>500</batchDuration>
+ <stopWhenIdle>N</stopWhenIdle>
+ <maxIdleTimeMs>500</maxIdleTimeMs>
+ <maxConsumeDurationMs>10000</maxConsumeDurationMs>
+
<directBootstrapServers>${BOOTSTRAP_SERVERS}</directBootstrapServers>
+ <AUTO_COMMIT>N</AUTO_COMMIT>
+ </transform>
+ """;
+ Node node = XmlHandler.loadXmlString(xml, "transform");
+ KafkaConsumerInputMeta meta =
+ XmlMetadataUtil.deSerializeFromXml(
+ node, KafkaConsumerInputMeta.class, new MemoryMetadataProvider());
+ assertEquals("10000", meta.getMaxConsumeDurationMs());
+ assertFalse(meta.isStopWhenIdle());
+ assertEquals("500", meta.getBatchDuration());
+ }
+
+ @Test
+ void testMaxConsumeDurationCheckZeroIsValid() {
+ KafkaConsumerInputMeta meta = new KafkaConsumerInputMeta();
+ meta.setMaxConsumeDurationMs("0");
+ List<ICheckResult> remarks = new ArrayList<>();
+ meta.check(
+ remarks,
+ null,
+ new TransformMeta(),
+ null,
+ new String[0],
+ new String[0],
+ null,
+ new Variables(),
+ null);
+ assertTrue(
+ remarks.stream()
+ .noneMatch(
+ r ->
+ r.getType() == ICheckResult.TYPE_RESULT_ERROR
+ && r.getText().contains("Max consume duration")));
}
/**
diff --git
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
new file mode 100644
index 0000000000..d2fb5e186e
--- /dev/null
+++
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
@@ -0,0 +1,60 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.pipeline.transforms.kafka.consumer;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import org.junit.jupiter.api.Test;
+
+class KafkaConsumerInputTest {
+
+ @Test
+ void maxConsumeDurationReachedWhenElapsed() {
+ assertFalse(KafkaConsumerInput.maxConsumeDurationReached(1_000L, 0L, 0L));
+ assertFalse(KafkaConsumerInput.maxConsumeDurationReached(1_000L, -5L, 0L));
+ assertFalse(KafkaConsumerInput.maxConsumeDurationReached(999L, 1_000L,
0L));
+ assertTrue(KafkaConsumerInput.maxConsumeDurationReached(1_000L, 1_000L,
0L));
+ assertTrue(KafkaConsumerInput.maxConsumeDurationReached(1_500L, 1_000L,
0L));
+ assertFalse(KafkaConsumerInput.maxConsumeDurationReached(1_500L, 1_000L,
600L));
+ }
+
+ @Test
+ void pollTimeoutUsesBatchDurationWhenNoDeadline() {
+ assertEquals(2_000L, KafkaConsumerInput.pollTimeoutMs(false, 2_000L, 0L,
0L, 0L));
+ assertEquals(Long.MAX_VALUE, KafkaConsumerInput.pollTimeoutMs(false, 0L,
0L, 0L, 0L));
+ }
+
+ @Test
+ void pollTimeoutUsesShortPollWhenStopWhenIdle() {
+ assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(true, 2_000L, 0L, 0L,
0L));
+ }
+
+ @Test
+ void pollTimeoutIsCappedToRemainingConsumeDuration() {
+ // A deadline uses a short poll so empty topics re-check the clock instead
of blocking in
+ // poll() for batchDuration (or forever when batchDuration is 0).
+ assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(false, 2_000L, 1_000L,
0L, 0L));
+ assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(false, 0L, 10_000L,
0L, 0L));
+ assertEquals(50L, KafkaConsumerInput.pollTimeoutMs(false, 2_000L, 1_000L,
0L, 950L));
+ assertEquals(0L, KafkaConsumerInput.pollTimeoutMs(false, 2_000L, 1_000L,
0L, 1_000L));
+ assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(true, 2_000L, 5_000L,
0L, 0L));
+ assertEquals(50L, KafkaConsumerInput.pollTimeoutMs(true, 2_000L, 1_000L,
0L, 950L));
+ }
+}