This is an automated email from the ASF dual-hosted git repository.

mattcasters 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 35dc7d3174 fix csv lazy conversion on Beam, fixes #3724 (#8551)
35dc7d3174 is described below

commit 35dc7d31744c830983f0cc5b28a087612ebec275
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Wed Sep 23 17:11:56 2026 +0200

    fix csv lazy conversion on Beam, fixes #3724 (#8551)
---
 .../0014-csv-lazy-conversion-validation.hpl        | 224 +++++++++++++
 .../beam_directrunner/0014-csv-lazy-conversion.hpl | 359 +++++++++++++++++++++
 .../datasets/0014-csv-lazy-conversion-golden.csv   |  61 ++++
 .../main-0014-csv-lazy-conversion.hwf              | 139 ++++++++
 .../dataset/0014-csv-lazy-conversion-golden.json   |  48 +++
 .../0014-csv-lazy-conversion-validation UNIT.json  |  44 +++
 .../apache/hop/beam/core/coder/HopRowCoder.java    |   6 +-
 .../core/transform/TransformBatchTransform.java    |   7 +-
 .../hop/beam/core/transform/TransformFn.java       |   4 +-
 .../org/apache/hop/beam/core/util/HopBeamUtil.java |  35 ++
 .../hop/beam/core/coder/HopRowCoderTest.java       |  40 +++
 .../apache/hop/beam/core/util/HopBeamUtilTest.java | 117 +++++++
 12 files changed, 1076 insertions(+), 8 deletions(-)

diff --git 
a/integration-tests/beam_directrunner/0014-csv-lazy-conversion-validation.hpl 
b/integration-tests/beam_directrunner/0014-csv-lazy-conversion-validation.hpl
new file mode 100644
index 0000000000..e3bf0ed5ae
--- /dev/null
+++ 
b/integration-tests/beam_directrunner/0014-csv-lazy-conversion-validation.hpl
@@ -0,0 +1,224 @@
+<?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>0014-csv-lazy-conversion-validation</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>2026/09/23 12:00:00.000</created_date>
+    <modified_user>-</modified_user>
+    <modified_date>2026/09/23 12:00:00.000</modified_date>
+  </info>
+  <notepads>
+  </notepads>
+  <order>
+    <hop>
+      <from>/tmp/0014/*.csv</from>
+      <to>Validate</to>
+      <enabled>Y</enabled>
+    </hop>
+  </order>
+  <transform>
+    <name>/tmp/0014/*.csv</name>
+    <type>TextFileInput2</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <accept_filenames>N</accept_filenames>
+    <passing_through_fields>N</passing_through_fields>
+    <accept_field/>
+    <accept_transform_name/>
+    <separator>,</separator>
+    <enclosure>"</enclosure>
+    <enclosure_breaks>N</enclosure_breaks>
+    <escapechar/>
+    <header>N</header>
+    <nr_headerlines>1</nr_headerlines>
+    <footer>N</footer>
+    <nr_footerlines>1</nr_footerlines>
+    <line_wrapped>N</line_wrapped>
+    <nr_wraps>1</nr_wraps>
+    <layout_paged>N</layout_paged>
+    <nr_lines_per_page>80</nr_lines_per_page>
+    <nr_lines_doc_header>0</nr_lines_doc_header>
+    <noempty>Y</noempty>
+    <include>N</include>
+    <include_field/>
+    <rownum>N</rownum>
+    <rownumByFile>N</rownumByFile>
+    <rownum_field/>
+    <format>Unix</format>
+    <encoding>UTF-8</encoding>
+    <length>Characters</length>
+    <add_to_result_filenames>Y</add_to_result_filenames>
+    <file>
+      <name>${java.io.tmpdir}/0014/</name>
+      <filemask>.*\.csv</filemask>
+      <exclude_filemask/>
+      <file_required>N</file_required>
+      <include_subfolders>N</include_subfolders>
+      <type>CSV</type>
+      <compression>None</compression>
+    </file>
+    <filters>
+    </filters>
+    <fields>
+      <field>
+        <name>stateCode</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <nullif/>
+        <ifnull/>
+        <position>-1</position>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+        <repeat>N</repeat>
+      </field>
+      <field>
+        <name>state</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <nullif/>
+        <ifnull/>
+        <position>-1</position>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+        <repeat>N</repeat>
+      </field>
+      <field>
+        <name>population</name>
+        <type>Integer</type>
+        <format>#</format>
+        <currency/>
+        <decimal/>
+        <group/>
+        <nullif/>
+        <ifnull/>
+        <position>-1</position>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+        <repeat>N</repeat>
+      </field>
+      <field>
+        <name>nrCustomers</name>
+        <type>Integer</type>
+        <format>#</format>
+        <currency/>
+        <decimal/>
+        <group/>
+        <nullif/>
+        <ifnull/>
+        <position>-1</position>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+        <repeat>N</repeat>
+      </field>
+      <field>
+        <name>sumIds</name>
+        <type>Integer</type>
+        <format>#</format>
+        <currency/>
+        <decimal/>
+        <group/>
+        <nullif/>
+        <ifnull/>
+        <position>-1</position>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+        <repeat>N</repeat>
+      </field>
+    </fields>
+    <limit>0</limit>
+    <error_ignored>N</error_ignored>
+    <skip_bad_files>N</skip_bad_files>
+    <file_error_field/>
+    <file_error_message_field/>
+    <error_line_skipped>N</error_line_skipped>
+    <error_count_field/>
+    <error_fields_field/>
+    <error_text_field/>
+    <bad_line_files_destination_directory/>
+    <bad_line_files_extension>warning</bad_line_files_extension>
+    <error_line_files_destination_directory/>
+    <error_line_files_extension>error</error_line_files_extension>
+    <line_number_files_destination_directory/>
+    <line_number_files_extension>line</line_number_files_extension>
+    <date_format_lenient>Y</date_format_lenient>
+    <date_format_locale>en_US</date_format_locale>
+    <shortFileFieldName/>
+    <pathFieldName/>
+    <hiddenFieldName/>
+    <lastModificationTimeFieldName/>
+    <uriNameFieldName/>
+    <rootUriNameFieldName/>
+    <extensionFieldName/>
+    <sizeFieldName/>
+    <attributes/>
+    <GUI>
+      <xloc>128</xloc>
+      <yloc>80</yloc>
+    </GUI>
+  </transform>
+  <transform>
+    <name>Validate</name>
+    <type>Dummy</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <attributes/>
+    <GUI>
+      <xloc>368</xloc>
+      <yloc>80</yloc>
+    </GUI>
+  </transform>
+  <transform_error_handling>
+  </transform_error_handling>
+  <attributes/>
+</pipeline>
diff --git a/integration-tests/beam_directrunner/0014-csv-lazy-conversion.hpl 
b/integration-tests/beam_directrunner/0014-csv-lazy-conversion.hpl
new file mode 100644
index 0000000000..d890f3a361
--- /dev/null
+++ b/integration-tests/beam_directrunner/0014-csv-lazy-conversion.hpl
@@ -0,0 +1,359 @@
+<?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>0014-csv-lazy-conversion</name>
+    <name_sync_with_filename>Y</name_sync_with_filename>
+    <description>CSV File Input with lazy conversion on the Beam engine: lazy 
values feed a Stream Lookup side input, its main input and a native Beam group 
by (#3724). state-population.txt mixes upper and lower case, so only the upper 
case states find a population.</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>2026/09/23 12:00:00.000</created_date>
+    <modified_user>-</modified_user>
+    <modified_date>2026/09/23 12:00:00.000</modified_date>
+  </info>
+  <notepads>
+  </notepads>
+  <order>
+    <hop>
+      <from>customers (lazy)</from>
+      <to>lookup population</to>
+      <enabled>Y</enabled>
+    </hop>
+    <hop>
+      <from>state population (lazy)</from>
+      <to>lookup population</to>
+      <enabled>Y</enabled>
+    </hop>
+    <hop>
+      <from>lookup population</from>
+      <to>per state</to>
+      <enabled>Y</enabled>
+    </hop>
+    <hop>
+      <from>per state</from>
+      <to>/tmp/0014/csv-lazy-conversion.csv</to>
+      <enabled>Y</enabled>
+    </hop>
+  </order>
+  <transform>
+    <name>customers (lazy)</name>
+    <type>CSVInput</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <filename>${PROJECT_HOME}/input/customers-noheader-1k.txt</filename>
+    <filename_field/>
+    <rownum_field/>
+    <include_filename>N</include_filename>
+    <separator>;</separator>
+    <enclosure/>
+    <header>N</header>
+    <buffer_size>50000</buffer_size>
+    <lazy_conversion>Y</lazy_conversion>
+    <add_filename_result>N</add_filename_result>
+    <parallel>N</parallel>
+    <newline_possible>N</newline_possible>
+    <encoding>UTF-8</encoding>
+    <fields>
+      <field>
+        <name>id</name>
+        <type>Integer</type>
+        <format>#</format>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>Last_name</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>First_name</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>cust_zip_code</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>city</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>birthdate</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>street</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>housenr</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>stateCode</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>state</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+    </fields>
+    <attributes/>
+    <GUI>
+      <xloc>112</xloc>
+      <yloc>96</yloc>
+    </GUI>
+  </transform>
+  <transform>
+    <name>state population (lazy)</name>
+    <type>CSVInput</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <filename>${PROJECT_HOME}/input/state-population.txt</filename>
+    <filename_field/>
+    <rownum_field/>
+    <include_filename>N</include_filename>
+    <separator>;</separator>
+    <enclosure/>
+    <header>N</header>
+    <buffer_size>50000</buffer_size>
+    <lazy_conversion>Y</lazy_conversion>
+    <add_filename_result>N</add_filename_result>
+    <parallel>N</parallel>
+    <newline_possible>N</newline_possible>
+    <encoding>UTF-8</encoding>
+    <fields>
+      <field>
+        <name>state</name>
+        <type>String</type>
+        <format/>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+      <field>
+        <name>population</name>
+        <type>Integer</type>
+        <format>#</format>
+        <currency/>
+        <decimal/>
+        <group/>
+        <length>-1</length>
+        <precision>-1</precision>
+        <trim_type>both</trim_type>
+      </field>
+    </fields>
+    <attributes/>
+    <GUI>
+      <xloc>336</xloc>
+      <yloc>224</yloc>
+    </GUI>
+  </transform>
+  <transform>
+    <name>lookup population</name>
+    <type>StreamLookup</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <from>state population (lazy)</from>
+    <input_sorted>N</input_sorted>
+    <preserve_memory>N</preserve_memory>
+    <sorted_list>N</sorted_list>
+    <integer_pair>N</integer_pair>
+    <lookup>
+      <key>
+        <name>state</name>
+        <field>state</field>
+      </key>
+      <value>
+        <name>population</name>
+        <rename>population</rename>
+        <default/>
+        <type>Integer</type>
+      </value>
+    </lookup>
+    <attributes/>
+    <GUI>
+      <xloc>336</xloc>
+      <yloc>96</yloc>
+    </GUI>
+  </transform>
+  <transform>
+    <name>per state</name>
+    <type>MemoryGroupBy</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <give_back_row>N</give_back_row>
+    <group>
+      <field>
+        <name>stateCode</name>
+      </field>
+      <field>
+        <name>state</name>
+      </field>
+      <field>
+        <name>population</name>
+      </field>
+    </group>
+    <fields>
+      <field>
+        <aggregate>nrCustomers</aggregate>
+        <subject>id</subject>
+        <type>COUNT_ALL</type>
+        <valuefield/>
+      </field>
+      <field>
+        <aggregate>sumIds</aggregate>
+        <subject>id</subject>
+        <type>SUM</type>
+        <valuefield/>
+      </field>
+    </fields>
+    <attributes/>
+    <GUI>
+      <xloc>544</xloc>
+      <yloc>96</yloc>
+    </GUI>
+  </transform>
+  <transform>
+    <name>/tmp/0014/csv-lazy-conversion.csv</name>
+    <type>BeamOutput</type>
+    <description/>
+    <distribute>Y</distribute>
+    <custom_distribution/>
+    <copies>1</copies>
+    <partitioning>
+      <method>none</method>
+      <schema_name/>
+    </partitioning>
+    <file_prefix>csv-lazy-conversion</file_prefix>
+    <file_suffix>.csv</file_suffix>
+    <output_location>${java.io.tmpdir}/0014/</output_location>
+    <windowed>N</windowed>
+    <attributes/>
+    <GUI>
+      <xloc>768</xloc>
+      <yloc>96</yloc>
+    </GUI>
+  </transform>
+  <transform_error_handling>
+  </transform_error_handling>
+  <attributes/>
+</pipeline>
diff --git 
a/integration-tests/beam_directrunner/datasets/0014-csv-lazy-conversion-golden.csv
 
b/integration-tests/beam_directrunner/datasets/0014-csv-lazy-conversion-golden.csv
new file mode 100644
index 0000000000..fa0b2c2632
--- /dev/null
+++ 
b/integration-tests/beam_directrunner/datasets/0014-csv-lazy-conversion-golden.csv
@@ -0,0 +1,61 @@
+stateCode,state,population,nrCustomers,sumIds
+AK,ALASKA,739795,19,9741
+AL,ALABAMA,,14,7303
+AR,ARKANSAS,,20,9160
+AS,AMERICAN SAMOA,,15,7163
+AZ,ARIZONA,,15,8194
+CA,CALIFORNIA,39536653,19,8189
+CO,COLORADO,,18,7096
+CT,CONNECTICUT,,14,7561
+DC,DISTRICT OF COLUMBIA,,16,6571
+DE,DELAWARE,,18,10869
+FL,FLORIDA,20984400,20,8398
+FM,FEDERATED STATES OF MICRONESIA,,12,6798
+GA,GEORGIA,,14,6183
+GU,GUAM,,19,11178
+HI,HAWAII,,22,10974
+IA,IOWA,,14,9589
+ID,IDAHO,,14,8842
+IL,ILLINOIS,,13,8218
+IN,INDIANA,6666818,20,8532
+KS,KANSAS,,18,8184
+KY,KENTUCKY,,18,8706
+LA,LOUISIANA,,21,9863
+MA,MASSACHUSETTS,,18,9569
+MD,MARYLAND,,18,9330
+ME,MAINE,,19,10207
+MH,MARSHALL ISLANDS,,14,7536
+MI,MICHIGAN,,11,6425
+MN,MINNESOTA,,23,10893
+MO,MISSOURI,,19,9353
+MP,NORTHERN MARIANA ISLANDS,,17,7677
+MS,MISSISSIPPI,,11,3870
+MT,MONTANA,,19,7693
+NC,NORTH CAROLINA,,16,5205
+ND,NORTH DAKOTA,,14,5510
+NE,NEBRASKA,1920076,20,10033
+NH,NEW HAMPSHIRE,,14,6125
+NJ,NEW JERSEY,,17,8315
+NM,NEW MEXICO,,10,3511
+NV,NEVADA,,24,11354
+NY,NEW YORK,19849399,21,10452
+OH,OHIO,,12,8861
+OK,OKLAHOMA,,26,12922
+OR,OREGON,,12,7451
+PA,PENNSYLVANIA,,15,7778
+PR,PUERTO RICO,,12,5730
+PW,PALAU,,16,8724
+RI,RHODE ISLAND,,20,11856
+SC,SOUTH CAROLINA,,8,4093
+SD,SOUTH DAKOTA,,19,8887
+TN,TENNESSEE,,11,8301
+TX,TEXAS,28304596,16,8215
+UT,UTAH,,19,9708
+VA,VIRGINIA,,20,9741
+VI,VIRGIN ISLANDS,,14,6275
+VT,VERMONT,,16,7575
+WA,WASHINGTON,7405743,18,9328
+WI,WISCONSIN,,17,10468
+WV,WEST VIRGINIA,,22,10315
+WY,WYOMING,,16,6990
+undefined,undefined,,13,6912
diff --git 
a/integration-tests/beam_directrunner/main-0014-csv-lazy-conversion.hwf 
b/integration-tests/beam_directrunner/main-0014-csv-lazy-conversion.hwf
new file mode 100644
index 0000000000..86293a2d65
--- /dev/null
+++ b/integration-tests/beam_directrunner/main-0014-csv-lazy-conversion.hwf
@@ -0,0 +1,139 @@
+<?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-0014-csv-lazy-conversion</name>
+  <name_sync_with_filename>Y</name_sync_with_filename>
+  <description/>
+  <extended_description/>
+  <workflow_version/>
+  <created_user>-</created_user>
+  <created_date>2026/09/23 12:00:00.000</created_date>
+  <modified_user>-</modified_user>
+  <modified_date>2026/09/23 12:00:00.000</modified_date>
+  <parameters>
+    </parameters>
+  <actions>
+    <action>
+      <name>Start</name>
+      <description/>
+      <type>SPECIAL</type>
+      <attributes/>
+      <repeat>N</repeat>
+      <schedulerType>0</schedulerType>
+      <intervalSeconds>0</intervalSeconds>
+      <intervalMinutes>60</intervalMinutes>
+      <hour>12</hour>
+      <minutes>0</minutes>
+      <weekDay>1</weekDay>
+      <DayOfMonth>1</DayOfMonth>
+      <parallel>N</parallel>
+      <xloc>96</xloc>
+      <yloc>96</yloc>
+      <attributes_hac/>
+    </action>
+    <action>
+      <name>delete /tmp/0014/*</name>
+      <description/>
+      <type>DELETE_FILES</type>
+      <attributes/>
+      <arg_from_previous>N</arg_from_previous>
+      <include_subfolders>N</include_subfolders>
+      <fields>
+        <field>
+          <name>${java.io.tmpdir}/0014/</name>
+          <filemask>.*</filemask>
+        </field>
+      </fields>
+      <parallel>N</parallel>
+      <xloc>256</xloc>
+      <yloc>96</yloc>
+      <attributes_hac/>
+    </action>
+    <action>
+      <name>0014-csv-lazy-conversion</name>
+      <description/>
+      <type>PIPELINE</type>
+      <attributes/>
+      <filename>${PROJECT_HOME}/0014-csv-lazy-conversion.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>
+      <set_logfile>N</set_logfile>
+      <logfile/>
+      <logext/>
+      <add_date>N</add_date>
+      <add_time>N</add_time>
+      <loglevel>Basic</loglevel>
+      <set_append_logfile>N</set_append_logfile>
+      <wait_until_finished>Y</wait_until_finished>
+      <create_parent_folder>N</create_parent_folder>
+      <run_configuration>local</run_configuration>
+      <parameters>
+        <pass_all_parameters>Y</pass_all_parameters>
+      </parameters>
+      <parallel>N</parallel>
+      <xloc>448</xloc>
+      <yloc>96</yloc>
+      <attributes_hac/>
+    </action>
+    <action>
+      <name>Run Pipeline Unit Tests</name>
+      <description/>
+      <type>RunPipelineTests</type>
+      <attributes/>
+      <test_names>
+        <test_name>
+          <name>0014-csv-lazy-conversion-validation UNIT</name>
+        </test_name>
+      </test_names>
+      <parallel>N</parallel>
+      <xloc>640</xloc>
+      <yloc>96</yloc>
+      <attributes_hac/>
+    </action>
+  </actions>
+  <hops>
+    <hop>
+      <from>Start</from>
+      <to>delete /tmp/0014/*</to>
+      <enabled>Y</enabled>
+      <evaluation>Y</evaluation>
+      <unconditional>Y</unconditional>
+    </hop>
+    <hop>
+      <from>delete /tmp/0014/*</from>
+      <to>0014-csv-lazy-conversion</to>
+      <enabled>Y</enabled>
+      <evaluation>Y</evaluation>
+      <unconditional>N</unconditional>
+    </hop>
+    <hop>
+      <from>0014-csv-lazy-conversion</from>
+      <to>Run Pipeline Unit Tests</to>
+      <enabled>Y</enabled>
+      <evaluation>Y</evaluation>
+      <unconditional>N</unconditional>
+    </hop>
+  </hops>
+  <notepads>
+  </notepads>
+  <attributes/>
+</workflow>
diff --git 
a/integration-tests/beam_directrunner/metadata/dataset/0014-csv-lazy-conversion-golden.json
 
b/integration-tests/beam_directrunner/metadata/dataset/0014-csv-lazy-conversion-golden.json
new file mode 100644
index 0000000000..c21422fe78
--- /dev/null
+++ 
b/integration-tests/beam_directrunner/metadata/dataset/0014-csv-lazy-conversion-golden.json
@@ -0,0 +1,48 @@
+{
+  "base_filename": "0014-csv-lazy-conversion-golden.csv",
+  "name": "0014-csv-lazy-conversion-golden",
+  "description": "",
+  "dataset_fields": [
+    {
+      "field_comment": "",
+      "field_length": -1,
+      "field_type": 2,
+      "field_precision": -1,
+      "field_name": "stateCode",
+      "field_format": ""
+    },
+    {
+      "field_comment": "",
+      "field_length": -1,
+      "field_type": 2,
+      "field_precision": -1,
+      "field_name": "state",
+      "field_format": ""
+    },
+    {
+      "field_comment": "",
+      "field_length": -1,
+      "field_type": 5,
+      "field_precision": 0,
+      "field_name": "population",
+      "field_format": "#"
+    },
+    {
+      "field_comment": "",
+      "field_length": -1,
+      "field_type": 5,
+      "field_precision": 0,
+      "field_name": "nrCustomers",
+      "field_format": "#"
+    },
+    {
+      "field_comment": "",
+      "field_length": -1,
+      "field_type": 5,
+      "field_precision": 0,
+      "field_name": "sumIds",
+      "field_format": "#"
+    }
+  ],
+  "folder_name": ""
+}
\ No newline at end of file
diff --git 
a/integration-tests/beam_directrunner/metadata/unit-test/0014-csv-lazy-conversion-validation
 UNIT.json 
b/integration-tests/beam_directrunner/metadata/unit-test/0014-csv-lazy-conversion-validation
 UNIT.json
new file mode 100644
index 0000000000..9f7ffece21
--- /dev/null
+++ 
b/integration-tests/beam_directrunner/metadata/unit-test/0014-csv-lazy-conversion-validation
 UNIT.json      
@@ -0,0 +1,44 @@
+{
+  "variableValues": [],
+  "database_replacements": [],
+  "autoOpening": true,
+  "basePath": "",
+  "golden_data_sets": [
+    {
+      "field_mappings": [
+        {
+          "transform_field": "stateCode",
+          "data_set_field": "stateCode"
+        },
+        {
+          "transform_field": "state",
+          "data_set_field": "state"
+        },
+        {
+          "transform_field": "population",
+          "data_set_field": "population"
+        },
+        {
+          "transform_field": "nrCustomers",
+          "data_set_field": "nrCustomers"
+        },
+        {
+          "transform_field": "sumIds",
+          "data_set_field": "sumIds"
+        }
+      ],
+      "field_order": [
+        "stateCode"
+      ],
+      "data_set_name": "0014-csv-lazy-conversion-golden",
+      "transform_name": "Validate"
+    }
+  ],
+  "input_data_sets": [],
+  "name": "0014-csv-lazy-conversion-validation UNIT",
+  "description": "",
+  "persist_filename": "",
+  "trans_test_tweaks": [],
+  "pipeline_filename": "./0014-csv-lazy-conversion-validation.hpl",
+  "test_type": "UNIT_TEST"
+}
\ No newline at end of file
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/coder/HopRowCoder.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/coder/HopRowCoder.java
index da4976440c..ac5244ff8f 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/coder/HopRowCoder.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/coder/HopRowCoder.java
@@ -154,7 +154,7 @@ public class HopRowCoder extends AtomicCoder<HopRow> {
       case IValueMeta.TYPE_BINARY:
         {
           byte[] bytes = (byte[]) object;
-          out.write(bytes.length);
+          out.writeInt(bytes.length);
           out.write(bytes);
         }
         break;
@@ -241,7 +241,7 @@ public class HopRowCoder extends AtomicCoder<HopRow> {
       case IValueMeta.TYPE_BINARY:
         {
           byte[] bytes = new byte[in.readInt()];
-          in.read(bytes);
+          in.readFully(bytes);
           return bytes;
         }
 
@@ -249,7 +249,7 @@ public class HopRowCoder extends AtomicCoder<HopRow> {
         {
           String hostname = (String) read(in, IValueMeta.TYPE_STRING);
           byte[] addr = new byte[in.readInt() == 1 ? 4 : 16];
-          in.read(addr);
+          in.readFully(addr);
           return InetAddress.getByAddress(hostname, addr);
         }
 
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
index feb9d1bc1d..4e133fd31c 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformBatchTransform.java
@@ -529,8 +529,9 @@ public class TransformBatchTransform extends 
TransformTransform {
             rowListener =
                 new RowAdapter() {
                   @Override
-                  public void rowWrittenEvent(IRowMeta rowMeta, Object[] row) {
-                    resultRows.add(row);
+                  public void rowWrittenEvent(IRowMeta rowMeta, Object[] row)
+                      throws HopTransformException {
+                    resultRows.add(HopBeamUtil.toNormalStorage(rowMeta, row));
                   }
                 };
             transformCombi.transform.addRowListener(rowListener);
@@ -566,7 +567,7 @@ public class TransformBatchTransform extends 
TransformTransform {
                       throws HopTransformException {
                     // We send the target row to a specific list...
                     //
-                    targetResultRows.add(row);
+                    targetResultRows.add(HopBeamUtil.toNormalStorage(rowMeta, 
row));
                   }
                 });
           }
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
index 756f85177d..543aacef20 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/transform/TransformFn.java
@@ -395,7 +395,7 @@ public class TransformFn extends TransformBaseFn {
             @Override
             public void rowWrittenEvent(IRowMeta rowMeta, Object[] row)
                 throws HopTransformException {
-              resultRows.add(new HopRow(row, rowMeta.size()));
+              resultRows.add(new HopRow(HopBeamUtil.toNormalStorage(rowMeta, 
row), rowMeta.size()));
             }
           };
       transformCombi.transform.addRowListener(rowListener);
@@ -427,7 +427,7 @@ public class TransformFn extends TransformBaseFn {
             public void rowReadEvent(IRowMeta rowMeta, Object[] row) throws 
HopTransformException {
               // We send the target row to a specific list...
               //
-              targetResultRows.add(row);
+              targetResultRows.add(HopBeamUtil.toNormalStorage(rowMeta, row));
             }
           });
     }
diff --git 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/util/HopBeamUtil.java
 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/util/HopBeamUtil.java
index b3e18322d3..21fa48aae5 100644
--- 
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/util/HopBeamUtil.java
+++ 
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/core/util/HopBeamUtil.java
@@ -19,7 +19,10 @@ package org.apache.hop.beam.core.util;
 
 import org.apache.hop.beam.core.HopRow;
 import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.exception.HopTransformException;
+import org.apache.hop.core.exception.HopValueException;
 import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.IValueMeta;
 import org.apache.hop.core.xml.XmlHandler;
 import org.apache.hop.core.xml.XmlHandlerCache;
 import org.apache.hop.metadata.api.IHopMetadataProvider;
@@ -62,6 +65,38 @@ public class HopBeamUtil {
     return new HopRow(newRow);
   }
 
+  /**
+   * Rows leaving a Hop transform in Beam are serialized with the HopRowCoder 
and described
+   * downstream by JSON row metadata, which only knows about normal storage. 
Values kept in lazy
+   * (binary string) or indexed storage, such as those of a CSV File Input 
with lazy conversion, are
+   * therefore converted to normal storage here. Values already in normal 
storage, including real
+   * binary data, are passed along untouched.
+   *
+   * @param rowMeta The row metadata of the transform output, including the 
storage information
+   * @param row The row to convert
+   * @return The given row if all values are in normal storage, a converted 
copy otherwise
+   * @throws HopTransformException In case a value can't be converted
+   */
+  public static Object[] toNormalStorage(IRowMeta rowMeta, Object[] row)
+      throws HopTransformException {
+    Object[] normalRow = row;
+    for (int i = 0; i < rowMeta.size(); i++) {
+      IValueMeta valueMeta = rowMeta.getValueMeta(i);
+      if (!valueMeta.isStorageNormal()) {
+        if (normalRow == row) {
+          normalRow = row.clone();
+        }
+        try {
+          normalRow[i] = valueMeta.convertToNormalStorageType(row[i]);
+        } catch (HopValueException e) {
+          throw new HopTransformException(
+              "Error converting field '" + valueMeta.getName() + "' to normal 
storage", e);
+        }
+      }
+    }
+    return normalRow;
+  }
+
   private static final Object object = new Object();
 
   public static void loadTransformMetadataFromXml(
diff --git 
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/coder/HopRowCoderTest.java
 
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/coder/HopRowCoderTest.java
index 0695d9afac..9a4f1593b9 100644
--- 
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/coder/HopRowCoderTest.java
+++ 
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/coder/HopRowCoderTest.java
@@ -17,12 +17,14 @@
 
 package org.apache.hop.beam.core.coder;
 
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
 import java.io.ByteArrayInputStream;
 import java.io.ByteArrayOutputStream;
 import java.io.IOException;
+import java.net.InetAddress;
 import java.sql.Timestamp;
 import java.util.Date;
 import org.apache.avro.Schema;
@@ -155,4 +157,42 @@ class HopRowCoderTest {
       assertEquals(genericRecord.get(key), verify.get(key));
     }
   }
+
+  @Test
+  void testEncodeDecodeBinary() throws IOException {
+    // Longer than 255 bytes so a length written as a single byte can't pass 
by accident
+    byte[] binary = new byte[300];
+    for (int i = 0; i < binary.length; i++) {
+      binary[i] = (byte) i;
+    }
+    HopRow row1 = new HopRow(new Object[] {"before", binary, new byte[0], 
"after", 42L});
+
+    hopRowCoder.encode(row1, outputStream);
+    HopRow row1d = hopRowCoder.decode(new 
ByteArrayInputStream(outputStream.toByteArray()));
+
+    Object[] decoded = row1d.getRow();
+    assertEquals(5, row1d.length());
+    assertEquals("before", decoded[0]);
+    assertArrayEquals(binary, (byte[]) decoded[1]);
+    assertArrayEquals(new byte[0], (byte[]) decoded[2]);
+    assertEquals("after", decoded[3]);
+    assertEquals(42L, decoded[4]);
+  }
+
+  @Test
+  void testEncodeDecodeInet() throws IOException {
+    InetAddress ipv4 = InetAddress.getByAddress("host4", new byte[] {10, 0, 0, 
1});
+    InetAddress ipv6 =
+        InetAddress.getByAddress(
+            "host6",
+            new byte[] {0x20, 0x01, 0x0d, (byte) 0xb8, 0, 0, 0, 0, 0, 0, 0, 0, 
0, 0, 0, 1});
+    HopRow row1 = new HopRow(new Object[] {ipv4, ipv6, "after"});
+
+    hopRowCoder.encode(row1, outputStream);
+    HopRow row1d = hopRowCoder.decode(new 
ByteArrayInputStream(outputStream.toByteArray()));
+
+    assertEquals(row1, row1d);
+    assertEquals("host4", ((InetAddress) row1d.getRow()[0]).getHostName());
+    assertEquals("host6", ((InetAddress) row1d.getRow()[1]).getHostName());
+  }
 }
diff --git 
a/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/util/HopBeamUtilTest.java
 
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/util/HopBeamUtilTest.java
new file mode 100644
index 0000000000..3a93bc122c
--- /dev/null
+++ 
b/plugins/engines/beam/src/test/java/org/apache/hop/beam/core/util/HopBeamUtilTest.java
@@ -0,0 +1,117 @@
+/*
+ * 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.beam.core.util;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.nio.charset.StandardCharsets;
+import org.apache.hop.core.exception.HopTransformException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.IValueMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaBinary;
+import org.apache.hop.core.row.value.ValueMetaInteger;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.junit.jupiter.api.Test;
+
+class HopBeamUtilTest {
+
+  /** Mimics a CSV File Input field with lazy conversion enabled. */
+  private static IValueMeta lazy(IValueMeta valueMeta) {
+    valueMeta.setStorageType(IValueMeta.STORAGE_TYPE_BINARY_STRING);
+    ValueMetaString storageMetadata = new ValueMetaString(valueMeta.getName());
+    storageMetadata.setConversionMask(valueMeta.getConversionMask());
+    storageMetadata.setTrimType(valueMeta.getTrimType());
+    valueMeta.setStorageMetadata(storageMetadata);
+    return valueMeta;
+  }
+
+  private static byte[] bytes(String string) {
+    return string.getBytes(StandardCharsets.UTF_8);
+  }
+
+  @Test
+  void testNormalStorageRowIsPassedAsIs() throws Exception {
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(new ValueMetaString("name"));
+    rowMeta.addValueMeta(new ValueMetaInteger("id"));
+    rowMeta.addValueMeta(new ValueMetaBinary("data"));
+
+    byte[] binary = {0, 1, 2, (byte) 255};
+    Object[] row = {"Hop", 42L, binary};
+
+    Object[] result = HopBeamUtil.toNormalStorage(rowMeta, row);
+
+    // Real binary data is normal storage: nothing is copied or converted
+    assertSame(row, result);
+    assertSame(binary, result[2]);
+  }
+
+  @Test
+  void testLazyValuesAreConvertedToNormalStorage() throws Exception {
+    ValueMetaInteger idMeta = new ValueMetaInteger("id");
+    idMeta.setConversionMask("#");
+    idMeta.setTrimType(IValueMeta.TRIM_TYPE_BOTH);
+
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(lazy(idMeta));
+    rowMeta.addValueMeta(lazy(new ValueMetaString("name")));
+    rowMeta.addValueMeta(lazy(new ValueMetaBinary("lazyBinary")));
+    rowMeta.addValueMeta(new ValueMetaBinary("realBinary"));
+    rowMeta.addValueMeta(lazy(new ValueMetaString("empty")));
+
+    byte[] realBinary = {0, 1, 2, (byte) 255};
+    Object[] row = {bytes(" 123"), bytes("Hop"), bytes("abc"), realBinary, 
null, "extra"};
+
+    Object[] result = HopBeamUtil.toNormalStorage(rowMeta, row);
+
+    assertEquals(123L, result[0]);
+    assertEquals("Hop", result[1]);
+    assertInstanceOf(byte[].class, result[2]);
+    assertArrayEquals(bytes("abc"), (byte[]) result[2]);
+    assertSame(realBinary, result[3]);
+    assertNull(result[4]);
+    // Over-allocated slots beyond the row metadata are kept
+    assertEquals(6, result.length);
+    assertEquals("extra", result[5]);
+
+    // The original row can still be in use by the local transform: it's not 
modified
+    assertNotSame(row, result);
+    assertArrayEquals(bytes(" 123"), (byte[]) row[0]);
+    assertArrayEquals(bytes("Hop"), (byte[]) row[1]);
+  }
+
+  @Test
+  void testConversionErrorNamesTheField() {
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(lazy(new ValueMetaInteger("quantity")));
+
+    HopTransformException exception =
+        assertThrows(
+            HopTransformException.class,
+            () -> HopBeamUtil.toNormalStorage(rowMeta, new Object[] 
{bytes("not a number")}));
+    assertTrue(exception.getMessage().contains("'quantity'"));
+  }
+}

Reply via email to