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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 2844adf2e8 [python] Support row-based blob upsert by key (#8435)
2844adf2e8 is described below

commit 2844adf2e8fd25c3509a075cbf9615cd8342118d
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Jul 2 19:58:11 2026 +0800

    [python] Support row-based blob upsert by key (#8435)
    
    Support row-based Python write/upsert APIs that can accept `Blob`
    objects without materializing them into bytes. This adds
    `TableWrite.write_row(...)` and `TableUpdate.upsert_by_key(rows,
    upsert_keys)` while preserving the existing Arrow upsert behavior.
---
 paimon-python/pypaimon/tests/blob_table_test.py    | 140 +++++++++++++
 .../pypaimon/tests/table_upsert_by_key_test.py     |  49 +++++
 .../pypaimon/tests/write/table_write_test.py       | 108 ++++++++++
 paimon-python/pypaimon/write/file_store_write.py   |  22 ++
 paimon-python/pypaimon/write/row_key_extractor.py  |  45 +++-
 paimon-python/pypaimon/write/row_utils.py          |  90 ++++++++
 paimon-python/pypaimon/write/table_update.py       |  30 +++
 .../pypaimon/write/table_update_by_row_id.py       | 139 ++++++++++++-
 .../pypaimon/write/table_upsert_by_key.py          | 230 +++++++++++++++++++++
 paimon-python/pypaimon/write/table_write.py        |  15 ++
 .../write/writer/dedicated_format_writer.py        |  84 +++++++-
 11 files changed, 942 insertions(+), 10 deletions(-)

diff --git a/paimon-python/pypaimon/tests/blob_table_test.py 
b/paimon-python/pypaimon/tests/blob_table_test.py
index 24b527c49b..e45b4c77bf 100755
--- a/paimon-python/pypaimon/tests/blob_table_test.py
+++ b/paimon-python/pypaimon/tests/blob_table_test.py
@@ -15,6 +15,7 @@
 # specific language governing permissions and limitations
 # under the License.
 
+import io
 import os
 import shutil
 import struct
@@ -26,9 +27,26 @@ import pyarrow as pa
 from pypaimon import CatalogFactory, Schema
 from pypaimon.schema.schema_change import SchemaChange
 from pypaimon.table.file_store_table import FileStoreTable
+from pypaimon.table.row.blob import Blob
 from pypaimon.write.commit_message import CommitMessage
 
 
+class _StreamingOnlyBlob(Blob):
+    def __init__(self, data: bytes):
+        self.data = data
+        self.opened = False
+
+    def to_data(self) -> bytes:
+        raise AssertionError("streaming blob should not be materialized")
+
+    def to_descriptor(self):
+        raise RuntimeError("streaming blob has no descriptor")
+
+    def new_input_stream(self):
+        self.opened = True
+        return io.BytesIO(self.data)
+
+
 class DedicatedFormatWriterTest(unittest.TestCase):
     """Tests for DedicatedFormatWriter functionality with paimon table 
operations."""
 
@@ -1059,6 +1077,128 @@ class DedicatedFormatWriterTest(unittest.TestCase):
             [b'first_blob', None, b'third_blob', None, b'fifth_blob'],
         )
 
+    def test_write_row_accepts_streaming_blob(self):
+        from pypaimon import Schema
+        from pypaimon.table.row.generic_row import GenericRow
+
+        pa_schema = pa.schema([
+            ('id', pa.int32()),
+            ('blob_data', pa.large_binary()),
+        ])
+        schema = Schema.from_pyarrow_schema(
+            pa_schema,
+            options={
+                'row-tracking.enabled': 'true',
+                'data-evolution.enabled': 'true'
+            }
+        )
+        self.catalog.create_table('test_db.blob_write_row_streaming', schema, 
False)
+        table = self.catalog.get_table('test_db.blob_write_row_streaming')
+
+        blob = _StreamingOnlyBlob(b'row-blob')
+        row = GenericRow([1, blob], table.fields)
+
+        write_builder = table.new_batch_write_builder()
+        writer = write_builder.new_write()
+        writer.write_row(row)
+        write_builder.new_commit().commit(writer.prepare_commit())
+        writer.close()
+
+        result = table.new_read_builder().new_read().to_arrow(
+            table.new_read_builder().new_scan().plan().splits())
+        self.assertTrue(blob.opened)
+        self.assertEqual(result.column('id').to_pylist(), [1])
+        self.assertEqual(result.column('blob_data').to_pylist(), [b'row-blob'])
+
+    def test_upsert_by_key_accepts_streaming_blob_row(self):
+        from pypaimon import Schema
+        from pypaimon.table.row.generic_row import GenericRow
+
+        pa_schema = pa.schema([
+            ('id', pa.int32()),
+            ('blob_data', pa.large_binary()),
+        ])
+        schema = Schema.from_pyarrow_schema(
+            pa_schema,
+            options={
+                'row-tracking.enabled': 'true',
+                'data-evolution.enabled': 'true'
+            }
+        )
+        self.catalog.create_table('test_db.blob_upsert_row_streaming', schema, 
False)
+        table = self.catalog.get_table('test_db.blob_upsert_row_streaming')
+
+        initial = pa.Table.from_pydict({
+            'id': [1],
+            'blob_data': [b'old-blob'],
+        }, schema=pa_schema)
+        write_builder = table.new_batch_write_builder()
+        writer = write_builder.new_write()
+        writer.write_arrow(initial)
+        write_builder.new_commit().commit(writer.prepare_commit())
+        writer.close()
+
+        updated_blob = _StreamingOnlyBlob(b'updated-blob')
+        update_row = GenericRow([1, updated_blob], table.fields)
+        shadow_blob = _StreamingOnlyBlob(b'shadow-blob')
+        shadow_row = GenericRow([2, shadow_blob], table.fields)
+        new_blob = _StreamingOnlyBlob(b'new-blob')
+        new_row = GenericRow([2, new_blob], table.fields)
+        upsert_builder = table.new_batch_write_builder()
+        table_update = 
upsert_builder.new_update().with_update_type(['blob_data'])
+        upsert_builder.new_commit().commit(
+            table_update.upsert_by_key(
+                [update_row, shadow_row, new_row], ['id']))
+
+        result = table.new_read_builder().new_read().to_arrow(
+            table.new_read_builder().new_scan().plan().splits())
+        rows = sorted(
+            result.select(['id', 'blob_data']).to_pylist(),
+            key=lambda item: item['id'],
+        )
+        self.assertTrue(updated_blob.opened)
+        self.assertFalse(shadow_blob.opened)
+        self.assertTrue(new_blob.opened)
+        self.assertEqual(
+            rows,
+            [
+                {'id': 1, 'blob_data': b'updated-blob'},
+                {'id': 2, 'blob_data': b'new-blob'},
+            ],
+        )
+
+    def test_upsert_by_key_rejects_heterogeneous_append_row_fields(self):
+        from pypaimon import Schema
+        from pypaimon.table.row.generic_row import GenericRow
+
+        pa_schema = pa.schema([
+            ('id', pa.int32()),
+            ('name', pa.string()),
+            ('blob_data', pa.large_binary()),
+        ])
+        schema = Schema.from_pyarrow_schema(
+            pa_schema,
+            options={
+                'row-tracking.enabled': 'true',
+                'data-evolution.enabled': 'true'
+            }
+        )
+        self.catalog.create_table(
+            'test_db.blob_upsert_row_heterogeneous_fields', schema, False)
+        table = self.catalog.get_table(
+            'test_db.blob_upsert_row_heterogeneous_fields')
+
+        partial_fields = [
+            table.field_dict['id'],
+            table.field_dict['blob_data'],
+        ]
+        partial_row = GenericRow([1, b'a'], partial_fields)
+        full_row = GenericRow([2, 'bob', b'b'], table.fields)
+
+        with self.assertRaisesRegex(ValueError, 'same field set'):
+            table.new_batch_write_builder().new_update().upsert_by_key(
+                [partial_row, full_row], ['id'])
+
     def test_update_blob_column(self):
         from pypaimon import Schema
         from pypaimon.read.reader.format_blob_reader import FormatBlobReader
diff --git a/paimon-python/pypaimon/tests/table_upsert_by_key_test.py 
b/paimon-python/pypaimon/tests/table_upsert_by_key_test.py
index cf612a968c..bd05a276ff 100644
--- a/paimon-python/pypaimon/tests/table_upsert_by_key_test.py
+++ b/paimon-python/pypaimon/tests/table_upsert_by_key_test.py
@@ -55,6 +55,9 @@ class _TableUpsertByKeyTestBase(DataEvolutionTestBase):
     def _apply_upsert(self, table_update, data, upsert_keys, cid):
         raise NotImplementedError
 
+    def _apply_upsert_rows(self, table_update, rows, upsert_keys, cid):
+        raise NotImplementedError
+
     # ------------------------------------------------------------------
     # Helpers built on the primitives
     # ------------------------------------------------------------------
@@ -72,6 +75,18 @@ class _TableUpsertByKeyTestBase(DataEvolutionTestBase):
         tc.close()
         return msgs
 
+    def _upsert_rows(self, table, rows, upsert_keys, update_cols=None):
+        wb = self._make_write_builder(table)
+        tu = wb.new_update()
+        if update_cols:
+            tu.with_update_type(update_cols)
+        cid = self._next_commit_id()
+        msgs = self._apply_upsert_rows(tu, rows, upsert_keys, cid)
+        tc = wb.new_commit()
+        self._apply_commit(tc, msgs, cid)
+        tc.close()
+        return msgs
+
     # ==================================================================
     # Basic upsert tests (non-partitioned)
     # ==================================================================
@@ -139,6 +154,34 @@ class _TableUpsertByKeyTestBase(DataEvolutionTestBase):
         )
         self.assertEqual([(1, 'Alice'), (2, 'Bob_new'), (3, 'Carol')], rows)
 
+    def test_row_upsert_mixed_update_and_append(self):
+        from pypaimon.table.row.generic_row import GenericRow
+
+        table = self._create_table()
+        self._write_arrow(table, pa.Table.from_pydict({
+            'id': [1, 2],
+            'name': ['Alice', 'Bob'],
+            'age': [25, 30],
+            'city': ['NYC', 'LA'],
+        }, schema=self.pa_schema))
+
+        rows = [
+            GenericRow([2, 'Bob_row', 31, 'LA2'], table.fields),
+            GenericRow([3, 'Carol', 35, 'Chicago'], table.fields),
+        ]
+        self._upsert_rows(table, rows, upsert_keys=['id'])
+
+        result = self._read_all(table)
+        actual = {
+            row['id']: (row['name'], row['age'], row['city'])
+            for row in result.to_pylist()
+        }
+        self.assertEqual({
+            1: ('Alice', 25, 'NYC'),
+            2: ('Bob_row', 31, 'LA2'),
+            3: ('Carol', 35, 'Chicago'),
+        }, actual)
+
     def test_upsert_for_existing_table_duplicate_keys(self):
         table = self._create_table()
         self._write_arrow(table, pa.Table.from_pydict({
@@ -790,11 +833,17 @@ class _BatchModeMixin(BatchModeMixin):
     def _apply_upsert(self, table_update, data, upsert_keys, cid):
         return table_update.upsert_by_arrow_with_key(data, upsert_keys)
 
+    def _apply_upsert_rows(self, table_update, rows, upsert_keys, cid):
+        return table_update.upsert_by_key(rows, upsert_keys)
+
 
 class _StreamModeMixin(StreamModeMixin):
     def _apply_upsert(self, table_update, data, upsert_keys, cid):
         return table_update.upsert_by_arrow_with_key(data, upsert_keys, cid)
 
+    def _apply_upsert_rows(self, table_update, rows, upsert_keys, cid):
+        return table_update.upsert_by_key(rows, upsert_keys, cid)
+
 
 # ======================================================================
 # Concrete test classes
diff --git a/paimon-python/pypaimon/tests/write/table_write_test.py 
b/paimon-python/pypaimon/tests/write/table_write_test.py
index c22ebcd52c..975b04f3eb 100644
--- a/paimon-python/pypaimon/tests/write/table_write_test.py
+++ b/paimon-python/pypaimon/tests/write/table_write_test.py
@@ -72,6 +72,23 @@ class TableWriteTest(unittest.TestCase):
     def tearDownClass(cls):
         shutil.rmtree(cls.tempdir, ignore_errors=True)
 
+    @staticmethod
+    def _commit_rows(table, rows):
+        write_builder = table.new_batch_write_builder()
+        table_write = write_builder.new_write()
+        table_commit = write_builder.new_commit()
+        for row in rows:
+            table_write.write_row(row)
+        table_commit.commit(table_write.prepare_commit())
+        table_write.close()
+        table_commit.close()
+
+    @staticmethod
+    def _read_sorted(table, sort_keys):
+        read_builder = table.new_read_builder()
+        return read_builder.new_read().to_arrow(
+            read_builder.new_scan().plan().splits()).sort_by(sort_keys)
+
     def test_write_snapshot(self):
         schema = Schema.from_pyarrow_schema(self.pa_schema, 
partition_keys=['dt'])
         self.catalog.create_table('default.test_write_snapshot', schema, False)
@@ -98,6 +115,97 @@ class TableWriteTest(unittest.TestCase):
         self.assertEqual(True, snapshot_json.__contains__("baseManifestList"))
         self.assertEqual(False, snapshot_json.__contains__("nextRowId"))
 
+    def test_write_row_append_only_partitioned_table(self):
+        from pypaimon.table.row.generic_row import GenericRow
+
+        schema = Schema.from_pyarrow_schema(
+            self.pa_schema, partition_keys=['dt'])
+        self.catalog.create_table(
+            'default.test_write_row_append_only_partitioned', schema, False)
+        table = self.catalog.get_table(
+            'default.test_write_row_append_only_partitioned')
+
+        reordered_fields = [
+            table.field_dict['dt'],
+            table.field_dict['behavior'],
+            table.field_dict['item_id'],
+            table.field_dict['user_id'],
+        ]
+        rows = [
+            GenericRow(['p1', 'a', 1001, 1], reordered_fields),
+            GenericRow(['p2', 'b', 1002, 2], reordered_fields),
+        ]
+        self._commit_rows(table, rows)
+
+        expected = pa.Table.from_pydict({
+            'user_id': [1, 2],
+            'item_id': [1001, 1002],
+            'behavior': ['a', 'b'],
+            'dt': ['p1', 'p2'],
+        }, schema=self.pa_schema)
+        actual = self._read_sorted(table, 'user_id')
+        self.assertEqual(expected, actual)
+
+    def test_write_row_fixed_bucket_primary_key_table(self):
+        from pypaimon.table.row.generic_row import GenericRow
+
+        schema = Schema.from_pyarrow_schema(
+            self.pk_pa_schema,
+            partition_keys=['dt'],
+            primary_keys=['user_id', 'dt'],
+            options={'bucket': '2'},
+        )
+        self.catalog.create_table(
+            'default.test_write_row_fixed_bucket_pk', schema, False)
+        table = self.catalog.get_table(
+            'default.test_write_row_fixed_bucket_pk')
+
+        rows = [
+            GenericRow([1, 1001, 'a', 'p1'], table.fields),
+            GenericRow([2, 1002, 'b', 'p2'], table.fields),
+        ]
+        self._commit_rows(table, rows)
+
+        expected = pa.Table.from_pydict({
+            'user_id': [1, 2],
+            'item_id': [1001, 1002],
+            'behavior': ['a', 'b'],
+            'dt': ['p1', 'p2'],
+        }, schema=self.pk_pa_schema)
+        sort_keys = [('user_id', 'ascending'), ('dt', 'ascending')]
+        self.assertEqual(
+            expected.sort_by(sort_keys), self._read_sorted(table, sort_keys))
+
+    def test_write_row_dynamic_bucket_primary_key_table(self):
+        from pypaimon.table.row.generic_row import GenericRow
+
+        schema = Schema.from_pyarrow_schema(
+            self.pk_pa_schema,
+            partition_keys=['dt'],
+            primary_keys=['user_id', 'dt'],
+            options={'bucket': '-1'},
+        )
+        self.catalog.create_table(
+            'default.test_write_row_dynamic_bucket_pk', schema, False)
+        table = self.catalog.get_table(
+            'default.test_write_row_dynamic_bucket_pk')
+
+        rows = [
+            GenericRow([1, 1001, 'a', 'p1'], table.fields),
+            GenericRow([2, 1002, 'b', 'p2'], table.fields),
+        ]
+        self._commit_rows(table, rows)
+
+        expected = pa.Table.from_pydict({
+            'user_id': [1, 2],
+            'item_id': [1001, 1002],
+            'behavior': ['a', 'b'],
+            'dt': ['p1', 'p2'],
+        }, schema=self.pk_pa_schema)
+        sort_keys = [('user_id', 'ascending'), ('dt', 'ascending')]
+        self.assertEqual(
+            expected.sort_by(sort_keys), self._read_sorted(table, sort_keys))
+
     def test_multi_prepare_commit_ao(self):
         schema = Schema.from_pyarrow_schema(self.pa_schema, 
partition_keys=['dt'])
         self.catalog.create_table('default.test_append_only_parquet', schema, 
False)
diff --git a/paimon-python/pypaimon/write/file_store_write.py 
b/paimon-python/pypaimon/write/file_store_write.py
index 58ceb1c83c..9aed2f36fb 100644
--- a/paimon-python/pypaimon/write/file_store_write.py
+++ b/paimon-python/pypaimon/write/file_store_write.py
@@ -26,6 +26,7 @@ logger = logging.getLogger(__name__)
 
 from pypaimon.common.options.core_options import CoreOptions
 from pypaimon.write.commit_message import CommitMessage
+from pypaimon.write.row_utils import row_values_to_arrow_table
 from pypaimon.write.writer.append_only_data_writer import AppendOnlyDataWriter
 from pypaimon.write.writer.dedicated_format_writer import DedicatedFormatWriter
 from pypaimon.write.writer.data_vector_writer import DataVectorWriter
@@ -65,6 +66,27 @@ class FileStoreWrite:
         writer = self.data_writers[key]
         writer.write(data)
 
+    def write_row(self, partition: Tuple, bucket: int, row, values_by_name: 
dict):
+        key = (partition, bucket)
+        if key not in self.data_writers:
+            self.data_writers[key] = self._create_data_writer(partition, 
bucket, self.options)
+        writer = self.data_writers[key]
+        if hasattr(writer, 'write_row'):
+            writer.write_row(row)
+            return
+
+        column_names = (
+            self.write_cols
+            if self.write_cols is not None
+            else list(self.table.field_names)
+        )
+        data = row_values_to_arrow_table(
+            values_by_name,
+            self.table.table_schema.fields,
+            column_names,
+        )
+        writer.write(data.to_batches()[0])
+
     def _create_data_writer(self, partition: Tuple, bucket: int, options: 
CoreOptions) -> DataWriter:
         def max_seq_number():
             return self._seq_number_stats(partition).get(bucket, 1)
diff --git a/paimon-python/pypaimon/write/row_key_extractor.py 
b/paimon-python/pypaimon/write/row_key_extractor.py
index dad93bf1ed..b74ec5f03c 100644
--- a/paimon-python/pypaimon/write/row_key_extractor.py
+++ b/paimon-python/pypaimon/write/row_key_extractor.py
@@ -19,7 +19,7 @@ import math
 import random
 import struct
 from abc import ABC, abstractmethod
-from typing import Dict, List, Tuple
+from typing import Any, Dict, List, Tuple
 
 import pyarrow as pa
 
@@ -90,6 +90,15 @@ class RowKeyExtractor(ABC):
         buckets = self._extract_buckets_batch(data)
         return partitions, buckets
 
+    def extract_partition_bucket_row(
+            self, values_by_name: Dict[str, Any]) -> Tuple[Tuple, int]:
+        partition = tuple(
+            values_by_name[self.table_schema.fields[i].name]
+            for i in self.partition_indices
+        )
+        bucket = self._extract_bucket_row(values_by_name)
+        return partition, bucket
+
     def _get_field_indices(self, field_names: List[str]) -> List[int]:
         if not field_names:
             return []
@@ -113,6 +122,10 @@ class RowKeyExtractor(ABC):
     def _extract_buckets_batch(self, table: pa.RecordBatch) -> List[int]:
         """Extract bucket numbers for all rows. Must be implemented by 
subclasses."""
 
+    @abstractmethod
+    def _extract_bucket_row(self, values_by_name: Dict[str, Any]) -> int:
+        """Extract bucket number for a single row."""
+
 
 class FixedBucketRowKeyExtractor(RowKeyExtractor):
     """Fixed bucket mode extractor with configurable number of buckets."""
@@ -141,6 +154,14 @@ class FixedBucketRowKeyExtractor(RowKeyExtractor):
             for row_idx in range(data.num_rows)
         ]
 
+    def _extract_bucket_row(self, values_by_name: Dict[str, Any]) -> int:
+        return _bucket_from_hash(
+            self._binary_row_hash_code(
+                tuple(values_by_name[name] for name in self.bucket_keys)
+            ),
+            self.num_buckets,
+        )
+
     def _binary_row_hash_code(self, row_values: Tuple) -> int:
         row = GenericRow(list(row_values), self._bucket_key_fields, 
RowKind.INSERT)
         serialized = GenericRowSerializer.to_bytes(row)
@@ -160,6 +181,9 @@ class UnawareBucketRowKeyExtractor(RowKeyExtractor):
     def _extract_buckets_batch(self, data: pa.RecordBatch) -> List[int]:
         return [0] * data.num_rows
 
+    def _extract_bucket_row(self, values_by_name: Dict[str, Any]) -> int:
+        return 0
+
 
 _SHORT_MAX_VALUE = 32767
 
@@ -311,6 +335,22 @@ class DynamicBucketRowKeyExtractor(RowKeyExtractor):
                 self._assigner.assign(partitions[row_idx], key_hash))
         return buckets
 
+    def _extract_bucket_row(self, values_by_name: Dict[str, Any]) -> int:
+        key_hash = _hash_bytes_by_words(
+            GenericRowSerializer.to_bytes(
+                GenericRow(
+                    [values_by_name[name] for name in self.bucket_keys],
+                    self._bucket_key_fields,
+                    RowKind.INSERT,
+                )
+            )[4:]
+        )
+        partition = tuple(
+            values_by_name[self.table_schema.fields[i].name]
+            for i in self.partition_indices
+        )
+        return self._assigner.assign(partition, key_hash)
+
 
 class PostponeBucketRowKeyExtractor(RowKeyExtractor):
     """Extractor for unaware bucket mode (bucket = -1, no primary keys)."""
@@ -323,3 +363,6 @@ class PostponeBucketRowKeyExtractor(RowKeyExtractor):
 
     def _extract_buckets_batch(self, data: pa.RecordBatch) -> List[int]:
         return [BucketMode.POSTPONE_BUCKET.value] * data.num_rows
+
+    def _extract_bucket_row(self, values_by_name: Dict[str, Any]) -> int:
+        return BucketMode.POSTPONE_BUCKET.value
diff --git a/paimon-python/pypaimon/write/row_utils.py 
b/paimon-python/pypaimon/write/row_utils.py
new file mode 100644
index 0000000000..6e372542a3
--- /dev/null
+++ b/paimon-python/pypaimon/write/row_utils.py
@@ -0,0 +1,90 @@
+# 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.
+
+from typing import Any, Dict, Iterable, List
+
+import pyarrow as pa
+
+from pypaimon.schema.data_types import DataField, PyarrowFieldParser
+from pypaimon.table.row.blob import Blob
+from pypaimon.table.row.internal_row import InternalRow
+from pypaimon.table.row.vector import Vector
+
+
+def is_blob_field(field: DataField) -> bool:
+    return getattr(field.type, 'type', None) == 'BLOB'
+
+
+def row_to_named_values(
+        row: InternalRow, table_fields: List[DataField]) -> Dict[str, Any]:
+    if hasattr(row, 'fields') and getattr(row, 'fields') is not None:
+        row_fields = getattr(row, 'fields')
+        return {
+            field.name: row.get_field(i)
+            for i, field in enumerate(row_fields)
+        }
+
+    if len(row) != len(table_fields):
+        raise ValueError(
+            "Rows without field metadata must have the same arity as the "
+            f"table schema: {len(row)} != {len(table_fields)}."
+        )
+    return {
+        field.name: row.get_field(i)
+        for i, field in enumerate(table_fields)
+    }
+
+
+def require_columns(
+        values_by_name: Dict[str, Any],
+        column_names: Iterable[str],
+        context: str) -> None:
+    missing = [
+        column_name
+        for column_name in column_names
+        if column_name not in values_by_name
+    ]
+    if missing:
+        raise ValueError(
+            f"{context} requires row field(s) {missing}, "
+            f"but row only contains {list(values_by_name.keys())}."
+        )
+
+
+def value_for_arrow(value: Any) -> Any:
+    if isinstance(value, Vector):
+        return value.to_list()
+    if isinstance(value, Blob):
+        raise ValueError(
+            "Blob values cannot be converted to Arrow without materializing "
+            "the stream. Use a Row-aware blob write path."
+        )
+    return value
+
+
+def row_values_to_arrow_table(
+        values_by_name: Dict[str, Any],
+        fields: List[DataField],
+        column_names: List[str]) -> pa.Table:
+    field_by_name = {field.name: field for field in fields}
+    selected_fields = [field_by_name[name] for name in column_names]
+    schema = PyarrowFieldParser.from_paimon_schema(selected_fields)
+    data = {
+        name: [value_for_arrow(values_by_name[name])]
+        for name in column_names
+    }
+    return pa.Table.from_pydict(data, schema=schema)
diff --git a/paimon-python/pypaimon/write/table_update.py 
b/paimon-python/pypaimon/write/table_update.py
index f01d0dad78..478211ca55 100644
--- a/paimon-python/pypaimon/write/table_update.py
+++ b/paimon-python/pypaimon/write/table_update.py
@@ -190,6 +190,16 @@ class TableUpdate:
             self.table, self.commit_user, commit_identifier
         ).upsert(table, upsert_keys, self.update_cols)
 
+    def _upsert_by_key(
+            self,
+            rows,
+            upsert_keys: List[str],
+            commit_identifier: int,
+    ) -> List[CommitMessage]:
+        return TableUpsertByKey(
+            self.table, self.commit_user, commit_identifier
+        ).upsert_rows(rows, upsert_keys, self.update_cols)
+
     def _merge_into(
             self,
             source: Any,
@@ -499,6 +509,14 @@ class BatchTableUpdate(TableUpdate):
             table, upsert_keys, BATCH_COMMIT_IDENTIFIER
         )
 
+    def upsert_by_key(
+            self, rows, upsert_keys: List[str]
+    ) -> List[CommitMessage]:
+        """Upsert rows into an append-only table by key columns."""
+        return self._upsert_by_key(
+            rows, upsert_keys, BATCH_COMMIT_IDENTIFIER
+        )
+
     def update_by_predicate(
             self,
             predicate: Optional[Predicate],
@@ -561,6 +579,18 @@ class StreamTableUpdate(TableUpdate):
             table, upsert_keys, commit_identifier
         )
 
+    def upsert_by_key(
+            self,
+            rows,
+            upsert_keys: List[str],
+            commit_identifier: int,
+    ) -> List[CommitMessage]:
+        """Upsert rows into an append-only table by key columns,
+        tagging the produced commit messages with ``commit_identifier``."""
+        return self._upsert_by_key(
+            rows, upsert_keys, commit_identifier
+        )
+
     def update_by_predicate(
             self,
             predicate: Optional[Predicate],
diff --git a/paimon-python/pypaimon/write/table_update_by_row_id.py 
b/paimon-python/pypaimon/write/table_update_by_row_id.py
index a44bee4e64..623668ad11 100644
--- a/paimon-python/pypaimon/write/table_update_by_row_id.py
+++ b/paimon-python/pypaimon/write/table_update_by_row_id.py
@@ -17,7 +17,7 @@
 
 import bisect
 from dataclasses import dataclass, field
-from typing import Dict, List, Optional, Tuple
+from typing import Any, Dict, List, Optional, Tuple
 
 import numpy as np
 import pyarrow as pa
@@ -26,13 +26,18 @@ import pyarrow.compute as pc
 from pypaimon.manifest.schema.data_file_meta import DataFileMeta
 from pypaimon.read.split import DataSplit
 from pypaimon.read.table_read import TableRead
-from pypaimon.schema.data_types import DataField
+from pypaimon.schema.data_types import DataField, PyarrowFieldParser
 from pypaimon.table.row.blob import Blob
 from pypaimon.table.row.generic_row import GenericRow
 from pypaimon.table.special_fields import SpecialFields
 from pypaimon.utils.range import Range
 from pypaimon.write.commit_message import CommitMessage
 from pypaimon.write.file_store_write import FileStoreWrite
+from pypaimon.write.row_utils import (
+    require_columns,
+    row_to_named_values,
+    value_for_arrow,
+)
 from pypaimon.write.writer.blob_writer import BlobWriter
 
 
@@ -187,6 +192,82 @@ class TableUpdateByRowId:
 
         return self.commit_messages
 
+    def update_row_columns(
+            self,
+            row,
+            row_ids: List[int],
+            column_names: List[str],
+    ) -> List[CommitMessage]:
+        return self.update_rows_columns([row], [row_ids], column_names)
+
+    def update_rows_columns(
+            self,
+            rows: List,
+            row_ids_by_row: List[List[int]],
+            column_names: List[str],
+    ) -> List[CommitMessage]:
+        if not column_names:
+            raise ValueError("column_names cannot be empty")
+        if len(rows) != len(row_ids_by_row):
+            raise ValueError(
+                "rows and row_ids_by_row must have the same length: "
+                f"{len(rows)} != {len(row_ids_by_row)}"
+            )
+
+        values_by_row = [
+            row_to_named_values(row, self.table.table_schema.fields)
+            for row in rows
+        ]
+        for values_by_name in values_by_row:
+            require_columns(values_by_name, column_names, 
"update_rows_columns")
+
+        row_entries = []
+        for values_by_name, row_ids in zip(values_by_row, row_ids_by_row):
+            row_entries.extend((row_id, values_by_name) for row_id in row_ids)
+
+        if not row_entries:
+            return []
+
+        row_entries.sort(key=lambda item: item[0])
+
+        for col_name in column_names:
+            if col_name not in self.table.field_names:
+                raise ValueError(f"Column {col_name} not found in table 
schema")
+
+        arrays = [
+            pa.array([row_id for row_id, _ in row_entries], type=pa.int64())
+        ]
+        fields = [pa.field(SpecialFields.ROW_ID.name, pa.int64())]
+        blob_object_columns: Dict[str, List[Any]] = {}
+
+        for col_name in column_names:
+            if self._is_blob_column(col_name):
+                blob_object_columns[col_name] = [
+                    values_by_name[col_name]
+                    for _, values_by_name in row_entries
+                ]
+                continue
+
+            table_field = self.table.field_dict[col_name]
+            arrow_field = PyarrowFieldParser.from_paimon_field(table_field)
+            arrays.append(
+                pa.array(
+                    [
+                        value_for_arrow(values_by_name[col_name])
+                        for _, values_by_name in row_entries
+                    ],
+                    type=arrow_field.type,
+                )
+            )
+            fields.append(arrow_field)
+
+        update_data = pa.Table.from_arrays(arrays, schema=pa.schema(fields))
+        data_with_first_row_id = self._calculate_first_row_id(update_data)
+        self._write_by_first_row_id(
+            data_with_first_row_id, column_names, blob_object_columns)
+
+        return self.commit_messages
+
     def _calculate_first_row_id(self, data: pa.Table) -> pa.Table:
         """Append ``_FIRST_ROW_ID`` to *data* by looking up each ``_ROW_ID``.
 
@@ -225,10 +306,15 @@ class TableUpdateByRowId:
             pa.array(first_row_id_values, type=pa.int64()),
         )
 
-    def _write_by_first_row_id(self, data: pa.Table, column_names: List[str]):
+    def _write_by_first_row_id(
+            self,
+            data: pa.Table,
+            column_names: List[str],
+            blob_object_columns: Optional[Dict[str, List[Any]]] = None):
         """Write data grouped by first_row_id."""
         first_row_id_array = data[self.FIRST_ROW_ID_COLUMN]
         unique_first_row_ids = pc.unique(first_row_id_array).to_pylist()
+        first_row_id_values = first_row_id_array.to_pylist()
 
         for first_row_id in unique_first_row_ids:
             entry = self._first_row_id_index.get(first_row_id)
@@ -237,7 +323,23 @@ class TableUpdateByRowId:
             split, _files = entry
 
             group_data = data.filter(pc.equal(first_row_id_array, 
first_row_id))
-            self._write_group(split.partition, first_row_id, group_data, 
column_names)
+            group_blob_object_columns = None
+            if blob_object_columns:
+                group_indices = [
+                    i for i, value in enumerate(first_row_id_values)
+                    if value == first_row_id
+                ]
+                group_blob_object_columns = {
+                    col_name: [values[i] for i in group_indices]
+                    for col_name, values in blob_object_columns.items()
+                }
+            self._write_group(
+                split.partition,
+                first_row_id,
+                group_data,
+                column_names,
+                group_blob_object_columns,
+            )
 
     def _read_original_file_data(self, first_row_id: int, column_names: 
List[str]) -> Optional[pa.Table]:
         """Read original file data for the given first_row_id.
@@ -281,6 +383,7 @@ class TableUpdateByRowId:
             update_data: pa.Table,
             column_names: List[str],
             first_row_id: int,
+            blob_object_columns: Optional[Dict[str, List[Any]]] = None,
     ) -> Tuple[Optional[pa.Table], Dict[str, List[object]]]:
         """Merge update data with original data, preserving row order.
 
@@ -323,6 +426,7 @@ class TableUpdateByRowId:
         update_by_col = {
             col_name: update_data[col_name].combine_chunks()
             for col_name in column_names
+            if col_name in update_data.column_names
         }
         update_positions = {
             int(relative_index.as_py()): idx
@@ -332,8 +436,17 @@ class TableUpdateByRowId:
         # non-empty group, so update_positions is non-empty here.
         blob_row_count = max(update_positions) + 1
         for col_name in column_names:
-            update_col = update_by_col[col_name]
             if self._is_blob_column(col_name):
+                if blob_object_columns and col_name in blob_object_columns:
+                    update_values = blob_object_columns[col_name]
+                    blob_columns[col_name] = [
+                        update_values[update_positions[i]]
+                        if i in update_positions
+                        else Blob.PLACE_HOLDER
+                        for i in range(blob_row_count)
+                    ]
+                    continue
+                update_col = update_by_col[col_name]
                 blob_columns[col_name] = [
                     update_col[update_positions[i]].as_py()
                     if i in update_positions
@@ -341,6 +454,7 @@ class TableUpdateByRowId:
                     for i in range(blob_row_count)
                 ]
                 continue
+            update_col = update_by_col[col_name]
             original_col = original_data[col_name].combine_chunks()
             if update_col.type != original_col.type:
                 update_col = self._coerce_column(
@@ -399,8 +513,13 @@ class TableUpdateByRowId:
                 return getattr(table_field.type, 'type', None) == 'BLOB'
         return False
 
-    def _write_group(self, partition: GenericRow, first_row_id: int,
-                     data: pa.Table, column_names: List[str]):
+    def _write_group(
+            self,
+            partition: GenericRow,
+            first_row_id: int,
+            data: pa.Table,
+            column_names: List[str],
+            blob_object_columns: Optional[Dict[str, List[Any]]] = None):
         """Write a group of data with the same first_row_id.
 
         Reads the original file data, merges in the update values, and
@@ -408,7 +527,11 @@ class TableUpdateByRowId:
         """
         original_data = self._read_original_file_data(first_row_id, 
column_names)
         merged_data, blob_columns = self._merge_update_with_original(
-            original_data, data, column_names, first_row_id,
+            original_data,
+            data,
+            column_names,
+            first_row_id,
+            blob_object_columns,
         )
 
         partition_tuple = tuple(partition.values)
diff --git a/paimon-python/pypaimon/write/table_upsert_by_key.py 
b/paimon-python/pypaimon/write/table_upsert_by_key.py
index 4faf02fea0..d59cf3e960 100644
--- a/paimon-python/pypaimon/write/table_upsert_by_key.py
+++ b/paimon-python/pypaimon/write/table_upsert_by_key.py
@@ -21,8 +21,11 @@ from typing import Any, Dict, List, Optional, Tuple
 import pyarrow as pa
 
 from pypaimon.read.table_read import TableRead
+from pypaimon.table.row.blob import Blob
+from pypaimon.table.row.internal_row import InternalRow
 from pypaimon.table.special_fields import SpecialFields
 from pypaimon.write.commit_message import CommitMessage
+from pypaimon.write.row_utils import require_columns, row_to_named_values
 from pypaimon.write.table_update_by_row_id import TableUpdateByRowId
 from pypaimon.write.table_write import StreamTableWrite
 
@@ -97,6 +100,144 @@ class TableUpsertByKey:
 
         return all_commit_messages
 
+    def upsert_rows(self, rows, upsert_keys: List[str],
+                    update_cols: Optional[List[str]] = None) -> 
List[CommitMessage]:
+        row_list = self._normalize_rows(rows)
+        if not row_list:
+            raise ValueError("rows must not be empty.")
+
+        row_items = [
+            (row, row_to_named_values(row, self.table.table_schema.fields))
+            for row in row_list
+        ]
+        for _, values_by_name in row_items:
+            self._validate_row_inputs(values_by_name, upsert_keys, update_cols)
+
+        if update_cols is None or len(update_cols) == 
len(self.table.field_names):
+            effective_update_cols = None
+        else:
+            effective_update_cols = update_cols
+
+        commit_messages: List[CommitMessage] = []
+        for partition_spec, partition_items in 
self._group_rows_by_partition(row_items):
+            commit_messages.extend(
+                self._upsert_row_partition(
+                    partition_items,
+                    upsert_keys,
+                    partition_spec,
+                    effective_update_cols,
+                )
+            )
+        return commit_messages
+
+    @staticmethod
+    def _normalize_rows(rows) -> List:
+        if isinstance(rows, InternalRow):
+            return [rows]
+        return list(rows)
+
+    def _group_rows_by_partition(
+            self,
+            row_items: List[Tuple[Any, Dict[str, Any]]],
+    ) -> List[Tuple[Dict[str, Any], List[Tuple[Any, Dict[str, Any]]]]]:
+        if not self.table.partition_keys:
+            return [({}, row_items)]
+
+        partition_to_items: Dict[Tuple[Any, ...], List[Tuple[Any, Dict[str, 
Any]]]] = {}
+        for item in row_items:
+            _, values_by_name = item
+            part_tuple = tuple(
+                values_by_name[key] for key in self.table.partition_keys
+            )
+            partition_to_items.setdefault(part_tuple, []).append(item)
+
+        return [
+            (dict(zip(self.table.partition_keys, part_tuple)), items)
+            for part_tuple, items in partition_to_items.items()
+        ]
+
+    def _upsert_row_partition(
+            self,
+            row_items: List[Tuple[Any, Dict[str, Any]]],
+            upsert_keys: List[str],
+            partition_spec: Dict[str, Any],
+            update_cols: Optional[List[str]],
+    ) -> List[CommitMessage]:
+        partition_key_set = set(self.table.partition_keys)
+        match_keys = [k for k in upsert_keys if k not in partition_key_set]
+        input_key_tuples = [
+            tuple(values_by_name[k] for k in match_keys)
+            for _, values_by_name in row_items
+        ]
+        for key_tuple in input_key_tuples:
+            for value in key_tuple:
+                if isinstance(value, Blob):
+                    raise ValueError("Blob values are not supported as upsert 
keys.")
+
+        row_items, input_key_tuples = self._dedup_row_items_last_write_wins(
+            row_items, input_key_tuples, partition_spec)
+
+        key_to_row_ids = self._build_key_to_row_ids_map(
+            match_keys, partition_spec, set(input_key_tuples)
+        )
+
+        matched_items: List[Tuple[Any, Dict[str, Any]]] = []
+        matched_row_ids: List[List[int]] = []
+        new_items: List[Tuple[Any, Dict[str, Any]]] = []
+        for item, key_tuple in zip(row_items, input_key_tuples):
+            if key_tuple in key_to_row_ids:
+                matched_items.append(item)
+                matched_row_ids.append(key_to_row_ids[key_tuple])
+            else:
+                new_items.append(item)
+
+        commit_messages: List[CommitMessage] = []
+        if matched_items:
+            cols_to_update = (
+                list(update_cols)
+                if update_cols is not None
+                else list(self.table.field_names)
+            )
+            for _, values_by_name in matched_items:
+                require_columns(values_by_name, cols_to_update, 
"upsert_by_key")
+            commit_messages.extend(TableUpdateByRowId(
+                self.table, self.commit_user, self.commit_identifier,
+            ).update_rows_columns(
+                [row for row, _ in matched_items],
+                matched_row_ids,
+                cols_to_update,
+            ))
+
+        if new_items:
+            commit_messages.extend(self._append_rows(new_items))
+
+        return commit_messages
+
+    @staticmethod
+    def _dedup_row_items_last_write_wins(
+            row_items: List[Tuple[Any, Dict[str, Any]]],
+            input_key_tuples: List[_KeyTuple],
+            partition_spec: Dict[str, Any],
+    ) -> Tuple[List[Tuple[Any, Dict[str, Any]]], List[_KeyTuple]]:
+        key_to_last_idx: Dict[_KeyTuple, int] = {}
+        for i, key_tuple in enumerate(input_key_tuples):
+            key_to_last_idx[key_tuple] = i
+
+        if len(key_to_last_idx) == len(input_key_tuples):
+            return row_items, input_key_tuples
+
+        original_count = len(input_key_tuples)
+        dedup_indices = sorted(key_to_last_idx.values())
+        logger.warning(
+            "Deduplicated input rows from %d to %d in partition %s "
+            "(kept last occurrence).",
+            original_count, len(dedup_indices), partition_spec,
+        )
+        return (
+            [row_items[i] for i in dedup_indices],
+            [input_key_tuples[i] for i in dedup_indices],
+        )
+
     # ------------------------------------------------------------------
     # Partition grouping
     # ------------------------------------------------------------------
@@ -283,6 +424,52 @@ class TableUpsertByKey:
         # that partition columns can be stripped first.  The same non-partition
         # key may legally appear in different partitions.
 
+    def _validate_row_inputs(
+            self,
+            values_by_name: Dict[str, Any],
+            upsert_keys: List[str],
+            update_cols: Optional[List[str]]):
+        if not self.table.options.data_evolution_enabled():
+            raise ValueError(
+                "upsert_by_key requires 'data-evolution.enabled' = 'true'."
+            )
+
+        if not self.table.options.row_tracking_enabled():
+            raise ValueError(
+                "upsert_by_key requires 'row-tracking.enabled' = 'true'."
+            )
+
+        if not upsert_keys:
+            raise ValueError("upsert_keys must not be empty.")
+
+        for key in upsert_keys:
+            if key not in self.table.field_names:
+                raise ValueError(
+                    f"upsert_key '{key}' is not in table schema fields: 
{self.table.field_names}"
+                )
+
+        unknown_fields = [
+            field_name
+            for field_name in values_by_name
+            if field_name not in self.table.field_names
+        ]
+        if unknown_fields:
+            raise ValueError(
+                f"upsert_by_key got row field(s) {unknown_fields} "
+                f"that are not in table schema fields: 
{self.table.field_names}"
+            )
+
+        require_columns(values_by_name, upsert_keys, "upsert_by_key")
+        require_columns(values_by_name, self.table.partition_keys, 
"upsert_by_key")
+
+        if update_cols is not None:
+            for col in update_cols:
+                if col not in self.table.field_names:
+                    raise ValueError(
+                        f"Column '{col}' in update_cols is not in table schema 
fields: "
+                        f"{self.table.field_names}"
+                    )
+
     def _build_key_to_row_ids_map(
             self,
             match_keys: List[str],
@@ -397,3 +584,46 @@ class TableUpsertByKey:
             return table_write.prepare_commit(self.commit_identifier)
         finally:
             table_write.close()
+
+    def _append_rows(
+            self,
+            row_items: List[Tuple[Any, Dict[str, Any]]],
+    ) -> List[CommitMessage]:
+        all_ordered_cols = self._append_row_column_names(row_items)
+
+        table_write = StreamTableWrite(self.table, self.commit_user)
+        try:
+            table_write.with_write_type(all_ordered_cols)
+            for row, _ in row_items:
+                table_write.write_row(row)
+            return table_write.prepare_commit(self.commit_identifier)
+        finally:
+            table_write.close()
+
+    def _append_row_column_names(
+            self,
+            row_items: List[Tuple[Any, Dict[str, Any]]],
+    ) -> List[str]:
+        first_field_names = set(row_items[0][1])
+        for _, values_by_name in row_items[1:]:
+            field_names = set(values_by_name)
+            if field_names == first_field_names:
+                continue
+
+            missing_fields = [
+                name for name in self.table.field_names
+                if name in first_field_names and name not in field_names
+            ]
+            extra_fields = [
+                name for name in self.table.field_names
+                if name in field_names and name not in first_field_names
+            ]
+            raise ValueError(
+                "upsert_by_key requires appended rows in the same batch to "
+                "have the same field set. Compared with the first appended "
+                f"row, missing fields: {missing_fields}; "
+                f"extra fields: {extra_fields}."
+            )
+        return [
+            name for name in self.table.field_names if name in 
first_field_names
+        ]
diff --git a/paimon-python/pypaimon/write/table_write.py 
b/paimon-python/pypaimon/write/table_write.py
index 91eafa7536..9d1c8fad51 100644
--- a/paimon-python/pypaimon/write/table_write.py
+++ b/paimon-python/pypaimon/write/table_write.py
@@ -23,6 +23,7 @@ import pyarrow as pa
 from pypaimon.schema.data_types import PyarrowFieldParser
 from pypaimon.snapshot.snapshot import BATCH_COMMIT_IDENTIFIER
 from pypaimon.table.row.blob import BlobConsumer
+from pypaimon.write.row_utils import require_columns, row_to_named_values
 from pypaimon.write.commit_message import CommitMessage
 from pypaimon.write.file_store_write import FileStoreWrite
 
@@ -59,6 +60,20 @@ class TableWrite:
             sub_table = pa.compute.take(data, indices_array)
             self.file_store_write.write(partition, bucket, sub_table)
 
+    def write_row(self, row):
+        values_by_name = row_to_named_values(row, 
self.table.table_schema.fields)
+        column_names = (
+            self.file_store_write.write_cols
+            if self.file_store_write.write_cols is not None
+            else list(self.table.field_names)
+        )
+        require_columns(values_by_name, column_names, "write_row")
+        require_columns(values_by_name, self.table.partition_keys, "write_row")
+        partition, bucket = (
+            self.row_key_extractor.extract_partition_bucket_row(values_by_name)
+        )
+        self.file_store_write.write_row(partition, bucket, row, values_by_name)
+
     def write_pandas(self, dataframe):
         pa_schema = 
PyarrowFieldParser.from_paimon_schema(self.table.table_schema.fields)
         record_batch = pa.RecordBatch.from_pandas(dataframe, schema=pa_schema)
diff --git a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py 
b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
index 4dc4fca2a5..6cfbb25629 100644
--- a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
+++ b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
@@ -25,9 +25,15 @@ from pypaimon.common.options.core_options import 
CoreOptions, ChangelogProducer
 from pypaimon.data.timestamp import Timestamp
 from pypaimon.manifest.schema.data_file_meta import DataFileMeta
 from pypaimon.manifest.schema.simple_stats import SimpleStats
-from pypaimon.schema.data_types import VectorType
+from pypaimon.schema.data_types import PyarrowFieldParser, VectorType
 from pypaimon.table.row.blob import BlobConsumer
 from pypaimon.table.row.generic_row import GenericRow
+from pypaimon.write.row_utils import (
+    require_columns,
+    row_to_named_values,
+    row_values_to_arrow_table,
+    value_for_arrow,
+)
 from pypaimon.write.writer.data_writer import DataWriter
 
 logger = logging.getLogger(__name__)
@@ -214,6 +220,82 @@ class DedicatedFormatWriter(DataWriter):
             self.abort()
             raise e
 
+    def write_row(self, row):
+        try:
+            values_by_name = row_to_named_values(
+                row, self.table.table_schema.fields)
+            required_columns = (
+                list(self.normal_column_names)
+                + list(self.blob_file_column_names)
+                + list(self.vector_write_columns)
+            )
+            require_columns(values_by_name, required_columns, "write_row")
+
+            if self.normal_column_names:
+                normal_values = dict(values_by_name)
+                for field_name in self.normal_column_names:
+                    normal_values[field_name] = (
+                        self._normal_row_value(field_name, 
normal_values[field_name])
+                    )
+                normal_data = row_values_to_arrow_table(
+                    normal_values,
+                    self.table.table_schema.fields,
+                    self.normal_column_names,
+                ).to_batches()[0]
+                processed_normal = self._process_normal_data(normal_data)
+                if processed_normal is not None:
+                    if self.pending_normal_data is None:
+                        self.pending_normal_data = processed_normal
+                    else:
+                        self.pending_normal_data = self._merge_normal_data(
+                            self.pending_normal_data, processed_normal)
+
+            for blob_column in self.blob_file_column_names:
+                arrow_type = PyarrowFieldParser.from_paimon_type(
+                    self.table.field_dict[blob_column].type)
+                self.blob_writers[blob_column].write_blob(
+                    values_by_name[blob_column], arrow_type)
+
+            if self.vector_writer is not None and self.vector_write_columns:
+                vector_data = row_values_to_arrow_table(
+                    values_by_name,
+                    self.table.table_schema.fields,
+                    self.vector_write_columns,
+                ).to_batches()[0]
+                self.vector_writer.write(vector_data)
+
+            self.record_count += 1
+            if self._should_roll_normal():
+                self._close_current_writers()
+
+        except Exception as e:
+            logger.error("Exception occurs when writing row. Cleaning up.", 
exc_info=e)
+            self.abort()
+            raise e
+
+    def _normal_row_value(self, field_name: str, value):
+        if field_name in self.blob_descriptor_fields and value is not None:
+            from pypaimon.table.row.blob import Blob
+
+            if isinstance(value, Blob):
+                try:
+                    return value.to_descriptor().serialize()
+                except Exception as e:
+                    raise ValueError(
+                        "blob-descriptor-field row values must be serialized "
+                        "BlobDescriptor bytes or a Blob with a descriptor."
+                    ) from e
+            return value
+
+        if field_name in self.blob_view_fields and value is not None:
+            from pypaimon.table.row.blob import BlobView
+
+            if isinstance(value, BlobView):
+                return value.view_struct.serialize()
+            return value
+
+        return value_for_arrow(value)
+
     def prepare_commit(self) -> List[DataFileMeta]:
         # Close any remaining data
         self._close_current_writers()

Reply via email to