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&lt;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));
+  }
+}

Reply via email to