Akash3121 commented on code in PR #10154: URL: https://github.com/apache/paimon/pull/10154#discussion_r4089335230
########## paimon-python/pypaimon/write/native_write.py: ########## @@ -0,0 +1,162 @@ +# 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. + +"""Optional Rust data writer behind PyPaimon's batch and stream builders.""" + +import pyarrow as pa + +from pypaimon.read.native_plan import native_method_available +from pypaimon.schema.arrow_schema import normalize_arrow_strings +from pypaimon.schema.data_types import PyarrowFieldParser, is_blob_file_field +from pypaimon.write.commit_message_serializer import deserialize_commit_message +from pypaimon.write.native_commit import create_native_write_table +from pypaimon.write.row_utils import row_to_named_values, row_values_to_arrow_table + + +def native_write_available() -> bool: + """Check every binding entry point used by the writer bridge.""" + return all(native_method_available(type_name, method) for type_name, method in ( + ('Table', 'from_resolved_schema'), + ('BatchWriteBuilder', '_with_commit_user'), + ('BatchWriteBuilder', 'with_overwrite'), + ('BatchTableWrite', 'write_arrow'), + ('BatchTableWrite', 'prepare_commit'), + ('StreamWriteBuilder', 'with_commit_user'), + ('StreamTableWrite', 'write_arrow'), + ('StreamTableWrite', 'prepare_commit'), + ('CommitMessage', 'serialize'), + )) + + +def create_native_write(table, commit_user, static_partition=None, stream=False): + """Return a native writer if the table can use the filesystem write path.""" + if (not native_write_available() + or table.options.data_evolution_enabled() + or table.options.file_format() != 'parquet' + or any(is_blob_file_field(field) for field in table.table_schema.fields)): + return None + native_table = create_native_write_table(table) + if native_table is None: + return None + if stream: + builder = native_table.new_stream_write_builder().with_commit_user(commit_user) + else: + builder = native_table.new_batch_write_builder()._with_commit_user(commit_user) + if static_partition is not None: + builder = builder.with_overwrite(static_partition) + return NativeTableWrite(table, commit_user, static_partition, stream, + builder.new_write()) + + +class NativeTableWrite: + """Use Rust for Arrow batches while retaining PyPaimon's commit-message API. + + Advanced Python writer methods switch to the Python writer before the first + native data write. Once Rust has written data, switching would split one + logical write across two writers, so it is rejected. + """ + + def __init__(self, table, commit_user, static_partition, stream, native_writer): + self.table = table + self.commit_user = commit_user + self.static_partition = static_partition + self.stream = stream + self._native_writer = native_writer + self._python_writer = None + self._written = False + + def _switch_to_python(self): + if self._python_writer is not None: + return self._python_writer + if self._written: + raise RuntimeError( + 'Cannot switch to the Python writer after native data was written') + self._native_writer.close() + self._native_writer = None + from pypaimon.write.table_write import BatchTableWrite, StreamTableWrite + if self.stream: + self._python_writer = StreamTableWrite(self.table, self.commit_user) + else: + self._python_writer = BatchTableWrite( + self.table, self.commit_user, self.static_partition) + return self._python_writer + + def __getattr__(self, name): + if name.startswith('_'): + raise AttributeError(name) + return getattr(self._switch_to_python(), name) + + def write_arrow(self, data): + if self._python_writer is not None: + return self._python_writer.write_arrow(data) + if isinstance(data, pa.RecordBatch): + return self.write_arrow_batch(data) + for batch in data.to_batches(): + self.write_arrow_batch(batch) + + def write_arrow_batch(self, data): + if self._python_writer is not None: + return self._python_writer.write_arrow_batch(data) + data = normalize_arrow_strings(data) Review Comment: `write.native.enabled` currently rejects an Arrow schema that the existing PyPaimon writer deliberately accepts. `TableWrite._validate_pyarrow_schema` treats `binary` and `fixed_size_binary` as compatible, and `table_write_test.py::test_validate_schema_allows_binary_family_for_write_cols protects that behavior. This bridge only normalizes strings before sending the batch to Rust. The Rust binding maps both Paimon `BINARY` and `VARBINARY` to Arrow `Binary` and then requires exact Arrow type equality, so a valid pa.binary(4) batch for a `pa.binary()` table fails only when native writing is enabled. Because ` _written` is set before the native call, this cannot safely fall back afterward either. Please normalize Python-accepted binary-family inputs to the native target schema before marking the writer as written, and add an end-to-end native test that writes a `pa.binary(4)` batch to a table created with `pa.binary()` and commits it su ccessfully. Failure scenario: An application successfully writing fixed-width binary arrays into a BYTES table enables `write.native.enabled=true` ; the identical write now raises before producing a commit message. Minimal fix: Before `_written = True` , validate against the PyPaimon schema and cast Python-supported top-level fixed-size binary fields to the Arrow type expected by the Rust writer. Do not catch and retry a native write error, because it may already have produced files. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
