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'"));
+ }
+}