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 162fae3946 [core] Support source table FileIO for BLOB descriptors 
(#8693)
162fae3946 is described below

commit 162fae39462a8a6698c65231d44bd7d0cd1a1e78
Author: LsomeYeah <[email protected]>
AuthorDate: Tue Jul 21 14:49:11 2026 +0800

    [core] Support source table FileIO for BLOB descriptors (#8693)
---
 docs/docs/primary-key-table/blob-storage.md        |  19 ++
 docs/generated/core_configuration.html             |   6 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  13 ++
 .../org/apache/paimon/utils/UriReaderFactory.java  |  45 ++++-
 .../apache/paimon/utils/UriReaderFactoryTest.java  |  26 +++
 .../paimon/table/BlobDescriptorReaderFactory.java  |  70 +++++++
 .../table/BlobDescriptorReaderFactoryTest.java     | 192 ++++++++++++++++++
 .../org/apache/paimon/flink/FlinkRowWrapper.java   |  30 ++-
 .../flink/sink/BlobDescriptorResolvingRow.java     | 217 +++++++++++++++++++++
 .../org/apache/paimon/flink/sink/FlinkSink.java    |  27 ++-
 .../apache/paimon/flink/sink/FlinkSinkBuilder.java |  83 +++++---
 .../apache/paimon/flink/sink/StoreSinkWrite.java   |  18 ++
 .../paimon/flink/sink/StoreSinkWriteImpl.java      |  18 +-
 .../org/apache/paimon/flink/BlobTableITCase.java   |  30 +++
 .../apache/paimon/flink/FlinkRowWrapperTest.java   |  26 +++
 .../flink/sink/BlobDescriptorResolvingRowTest.java |  74 +++++++
 .../paimon/spark/SparkInternalRowWrapper.java      |  28 ++-
 .../java/org/apache/paimon/spark/SparkRow.java     |  18 +-
 .../write/PaimonV2MetadataAwareDataWriter.java     |   6 +-
 .../spark/commands/DataEvolutionPaimonWriter.scala |   3 +-
 .../paimon/spark/commands/PaimonSparkWriter.scala  |  26 ++-
 .../procedure/SparkPostponeCompactProcedure.scala  |  26 ++-
 .../apache/paimon/spark/util/SparkRowUtils.scala   |  16 +-
 .../spark/write/DataEvolutionTableDataWrite.scala  |   6 +-
 .../paimon/spark/write/PaimonBatchWriteBase.scala  |  30 +--
 .../paimon/spark/write/PaimonDataWrite.scala       |   6 +-
 .../paimon/spark/write/PaimonDeltaWriteBase.scala  |   2 +-
 .../paimon/spark/write/PaimonDeltaWriter.scala     |  10 +-
 .../paimon/spark/write/PaimonV2DataWriter.scala    |   7 +-
 .../apache/paimon/spark/write/WriteHelper.scala    |  11 +-
 paimon-spark/paimon-spark-ut/pom.xml               |   7 +
 .../apache/paimon/spark/SparkInternalRowTest.java  |  30 +++
 .../paimon/spark/write/PaimonBatchWriteTest.java   |  82 ++++++++
 .../org/apache/paimon/spark/sql/BlobTestBase.scala |  30 +++
 .../paimon/spark/sql/PostponeBucketTableTest.scala |  63 ++++++
 35 files changed, 1182 insertions(+), 119 deletions(-)

diff --git a/docs/docs/primary-key-table/blob-storage.md 
b/docs/docs/primary-key-table/blob-storage.md
index 938129c3f2..bbc9329e79 100644
--- a/docs/docs/primary-key-table/blob-storage.md
+++ b/docs/docs/primary-key-table/blob-storage.md
@@ -67,6 +67,25 @@ Reads return the payload bytes by default; the existing 
`blob-as-descriptor` rea
 A `blob-descriptor-field` is written inline to the normal data file and does 
not participate in managed storage or its
 reference sidecars.
 
+When descriptor-backed BLOBs are copied to another table, the target normally 
rebuilds a `FileIO` from its catalog
+context. For a managed `blob-field`, this is a copy flow: the target writes 
the payload into its own BLOB storage and
+does not retain the source descriptor. If the source table uses table-scoped 
credentials, configure
+`blob-descriptor.source-table` on the target so that the source table's 
`FileIO` is used to materialize the payload:
+
+```sql
+ALTER TABLE media_copy SET TBLPROPERTIES (
+    'blob-descriptor.source-table' = 'db.media$branch_rt'
+);
+```
+
+Use `blob-descriptor-field` to retain literal descriptors, or 
`blob-view-field` to retain a logical, no-copy reference
+to an upstream row. Other `blob-descriptor.*` filesystem options remain 
sufficient when the source storage can be
+accessed with static configuration; `source-table` is for table-scoped 
`FileIO` credentials.
+
+The source table must belong to the same catalog. A branch suffix is 
supported. Target tables without a catalog loader,
+including external tables in REST catalogs, are not supported. When this 
option is set, it takes precedence over other
+`blob-descriptor.*` options; remove it before switching back to 
descriptor-specific filesystem configuration.
+
 `ARRAY<BLOB>` is externalized element by element. Every non-null `Blob` 
element is copied into managed storage, while
 array order, a null array, and null elements are preserved. An empty array 
writes no payload. `ARRAY<BLOB>` uses
 `blob-field`; `blob-descriptor-field` and `blob-view-field` remain scalar-only 
declarations.
diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index bd362a0335..6c077fa9e4 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -74,6 +74,12 @@ under the License.
             <td>String</td>
             <td>Comma-separated field names to treat as BLOB fields and store 
as serialized BlobDescriptor bytes inline in data files.</td>
         </tr>
+        <tr>
+            <td><h5>blob-descriptor.source-table</h5></td>
+            <td style="word-wrap: break-word;">(none)</td>
+            <td>String</td>
+            <td>The source table whose FileIO is used to read 
descriptor-backed BLOB content and copy it into the target table's managed BLOB 
storage. The table must belong to the current catalog and can include a branch 
suffix, for example db.table$branch_rt. This option is not supported for target 
tables without a catalog loader, including external tables in REST catalogs. 
When set, other blob-descriptor.* FileIO options are ignored.</td>
+        </tr>
         <tr>
             <td><h5>blob-field</h5></td>
             <td style="word-wrap: break-word;">(none)</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 1419b7c2b2..6f6869eb18 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -2590,6 +2590,19 @@ public class CoreOptions implements Serializable {
                     .withDescription(
                             "Write blob field using blob descriptor rather 
than blob bytes.");
 
+    public static final ConfigOption<String> BLOB_DESCRIPTOR_SOURCE_TABLE =
+            key(BLOB_DESCRIPTOR_PREFIX + "source-table")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "The source table whose FileIO is used to read 
descriptor-backed BLOB "
+                                    + "content and copy it into the target 
table's managed BLOB "
+                                    + "storage. The table must belong to the 
current catalog and can "
+                                    + "include a branch suffix, for example 
db.table$branch_rt. This "
+                                    + "option is not supported for target 
tables without a catalog "
+                                    + "loader, including external tables in 
REST catalogs. When set, "
+                                    + "other blob-descriptor.* FileIO options 
are ignored.");
+
     public static final ConfigOption<Boolean> BLOB_WRITE_NULL_ON_MISSING_FILE =
             key("blob-write-null-on-missing-file")
                     .booleanType()
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/utils/UriReaderFactory.java 
b/paimon-common/src/main/java/org/apache/paimon/utils/UriReaderFactory.java
index 08f0691303..215a96295d 100644
--- a/paimon-common/src/main/java/org/apache/paimon/utils/UriReaderFactory.java
+++ b/paimon-common/src/main/java/org/apache/paimon/utils/UriReaderFactory.java
@@ -35,18 +35,25 @@ import java.util.concurrent.ConcurrentHashMap;
 /** A factory to create and cache {@link UriReader}. */
 public class UriReaderFactory implements Serializable {
 
-    private final CatalogContext context;
+    private static final long serialVersionUID = -8477284718943635074L;
+
+    @Nullable private final CatalogContext context;
     private transient Map<UriKey, UriReader> readers;
 
-    public UriReaderFactory(CatalogContext context) {
+    public UriReaderFactory(@Nullable CatalogContext context) {
         this.context = context;
         this.readers = new ConcurrentHashMap<>();
     }
 
+    /** Creates a factory which uses the provided {@link FileIO} for non-HTTP 
URIs. */
+    public static UriReaderFactory fromFileIO(FileIO fileIO) {
+        return new ProvidedFileIOUriReaderFactory(fileIO);
+    }
+
     public UriReader create(String input) {
         URI uri = parseUri(input);
         UriKey key = new UriKey(uri.getScheme(), uri.getAuthority());
-        return readers.computeIfAbsent(key, k -> newReader(k, uri));
+        return readers.computeIfAbsent(key, k -> newReader(uri));
     }
 
     private static URI parseUri(String input) {
@@ -73,19 +80,43 @@ public class UriReaderFactory implements Serializable {
         this.readers = new ConcurrentHashMap<>();
     }
 
-    private UriReader newReader(UriKey key, URI uri) {
-        if ("http".equals(key.scheme) || "https".equals(key.scheme)) {
+    protected UriReader newReader(URI uri) {
+        if (isHttp(uri)) {
             return UriReader.fromHttp();
         }
 
         try {
-            FileIO fileIO = FileIO.get(new Path(uri), context);
-            return UriReader.fromFile(fileIO);
+            FileIO createdFileIO = FileIO.get(new Path(uri), 
Objects.requireNonNull(context));
+            return UriReader.fromFile(createdFileIO);
         } catch (IOException e) {
             throw new RuntimeException(e);
         }
     }
 
+    private static boolean isHttp(URI uri) {
+        return "http".equals(uri.getScheme()) || 
"https".equals(uri.getScheme());
+    }
+
+    private static final class ProvidedFileIOUriReaderFactory extends 
UriReaderFactory {
+
+        private static final long serialVersionUID = 1L;
+
+        // Intentionally not transient. FileIO is serializable by contract, 
while implementations
+        // keep process-local clients transient. Distributed workers need this 
serialized FileIO to
+        // rebuild the transient reader cache with table-scoped credentials.
+        private final FileIO fileIO;
+
+        private ProvidedFileIOUriReaderFactory(FileIO fileIO) {
+            super(null);
+            this.fileIO = Objects.requireNonNull(fileIO);
+        }
+
+        @Override
+        protected UriReader newReader(URI uri) {
+            return isHttp(uri) ? super.newReader(uri) : 
UriReader.fromFile(fileIO);
+        }
+    }
+
     private static final class UriKey {
 
         private final @Nullable String scheme;
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/utils/UriReaderFactoryTest.java 
b/paimon-common/src/test/java/org/apache/paimon/utils/UriReaderFactoryTest.java
index e0679cb74f..5e1fb018ac 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/utils/UriReaderFactoryTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/utils/UriReaderFactoryTest.java
@@ -19,6 +19,9 @@
 package org.apache.paimon.utils;
 
 import org.apache.paimon.catalog.CatalogContext;
+import org.apache.paimon.fs.IsolatedDirectoryFileIO;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.SeekableInputStream;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.utils.UriReader.FileUriReader;
 import org.apache.paimon.utils.UriReader.HttpUriReader;
@@ -94,6 +97,29 @@ public class UriReaderFactoryTest {
         assertThat(reader).isInstanceOf(FileUriReader.class);
     }
 
+    @Test
+    public void testProvidedFileIOSurvivesSerialization() throws Exception {
+        java.nio.file.Path file = tempPath.resolve("file.txt");
+        Files.write(file, new byte[] {1, 2});
+
+        Options options = new Options();
+        options.set(IsolatedDirectoryFileIO.ROOT_DIR, new 
Path(tempPath.toUri()).toString());
+        IsolatedDirectoryFileIO fileIO = new IsolatedDirectoryFileIO();
+        fileIO.configure(CatalogContext.create(options));
+
+        String fileUri = file.toUri().toString();
+        UriReaderFactory originalFactory = UriReaderFactory.fromFileIO(fileIO);
+        UriReader originalReader = originalFactory.create(fileUri);
+        UriReaderFactory fileIOFactory = 
InstantiationUtil.clone(originalFactory);
+        UriReader deserializedReader = fileIOFactory.create(fileUri);
+
+        assertThat(deserializedReader).isNotSameAs(originalReader);
+        try (SeekableInputStream inputStream = 
deserializedReader.newInputStream(fileUri)) {
+            assertThat(inputStream.read()).isEqualTo(1);
+            assertThat(inputStream.read()).isEqualTo(2);
+        }
+    }
+
     @Test
     public void testCreateUriReaderWithAuthority() {
         UriReader reader1 = 
factory.create("http://my_bucket1/path/to/file.txt";);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/BlobDescriptorReaderFactory.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/BlobDescriptorReaderFactory.java
new file mode 100644
index 0000000000..a7c29ebfdb
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/BlobDescriptorReaderFactory.java
@@ -0,0 +1,70 @@
+/*
+ * 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.paimon.table;
+
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.CatalogContext;
+import org.apache.paimon.catalog.CatalogLoader;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.utils.BlobDescriptorUtils;
+import org.apache.paimon.utils.UriReaderFactory;
+
+import static org.apache.paimon.CoreOptions.BLOB_DESCRIPTOR_SOURCE_TABLE;
+import static org.apache.paimon.utils.Preconditions.checkNotNull;
+
+/** Creates the {@link UriReaderFactory} used to resolve BLOB descriptors 
during writes. */
+public final class BlobDescriptorReaderFactory {
+
+    private BlobDescriptorReaderFactory() {}
+
+    public static UriReaderFactory create(FileStoreTable table) {
+        Options tableOptions = table.coreOptions().toConfiguration();
+        String sourceTable = tableOptions.get(BLOB_DESCRIPTOR_SOURCE_TABLE);
+        if (sourceTable != null) {
+            return fromSourceTable(table, sourceTable);
+        }
+
+        CatalogContext descriptorContext =
+                BlobDescriptorUtils.getCatalogContext(
+                        table.catalogEnvironment().catalogContext(), 
tableOptions);
+        return new UriReaderFactory(descriptorContext);
+    }
+
+    private static UriReaderFactory fromSourceTable(FileStoreTable table, 
String sourceTable) {
+        CatalogLoader catalogLoader =
+                checkNotNull(
+                        table.catalogEnvironment().catalogLoader(),
+                        "Option '%s' is not supported for tables without a 
catalog loader, "
+                                + "including external tables in REST 
catalogs.",
+                        BLOB_DESCRIPTOR_SOURCE_TABLE.key());
+        Identifier sourceIdentifier = Identifier.fromString(sourceTable);
+        try (Catalog catalog = catalogLoader.load()) {
+            FileIO sourceFileIO = catalog.getTable(sourceIdentifier).fileIO();
+            // Initialize lazy credentials before serializing FileIO to 
distributed workers.
+            sourceFileIO.isObjectStore();
+            return UriReaderFactory.fromFileIO(sourceFileIO);
+        } catch (Exception e) {
+            throw new RuntimeException(
+                    String.format("Failed to load BLOB descriptor source table 
'%s'.", sourceTable),
+                    e);
+        }
+    }
+}
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/BlobDescriptorReaderFactoryTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/BlobDescriptorReaderFactoryTest.java
new file mode 100644
index 0000000000..7aaf3d5b18
--- /dev/null
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/BlobDescriptorReaderFactoryTest.java
@@ -0,0 +1,192 @@
+/*
+ * 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.paimon.table;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.CatalogContext;
+import org.apache.paimon.catalog.CatalogLoader;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.IsolatedDirectoryFileIO;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.SeekableInputStream;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.rest.RESTApi;
+import org.apache.paimon.rest.RESTTokenFileIO;
+import org.apache.paimon.rest.responses.GetTableTokenResponse;
+import org.apache.paimon.utils.InstantiationUtil;
+import org.apache.paimon.utils.UriReaderFactory;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Files;
+import java.util.Collections;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests for {@link BlobDescriptorReaderFactory}. */
+public class BlobDescriptorReaderFactoryTest {
+
+    @TempDir java.nio.file.Path tempPath;
+
+    @Test
+    public void testUseCatalogContextByDefault() throws Exception {
+        java.nio.file.Path tableDirectory = 
Files.createDirectory(tempPath.resolve("table"));
+        java.nio.file.Path blobFile = tableDirectory.resolve("blob");
+        Files.write(blobFile, new byte[] {1, 2});
+
+        Options catalogOptions = new Options();
+        catalogOptions.set(
+                IsolatedDirectoryFileIO.ROOT_DIR, "isolated://" + 
tableDirectory.toString());
+        CatalogEnvironment catalogEnvironment = mock(CatalogEnvironment.class);
+        
when(catalogEnvironment.catalogContext()).thenReturn(CatalogContext.create(catalogOptions));
+        FileStoreTable table = mock(FileStoreTable.class);
+        when(table.catalogEnvironment()).thenReturn(catalogEnvironment);
+        
when(table.coreOptions()).thenReturn(CoreOptions.fromMap(Collections.emptyMap()));
+
+        UriReaderFactory readerFactory = 
BlobDescriptorReaderFactory.create(table);
+        String blobUri = "isolated://" + blobFile;
+        try (SeekableInputStream inputStream =
+                readerFactory.create(blobUri).newInputStream(blobUri)) {
+            assertThat(inputStream.read()).isEqualTo(1);
+            assertThat(inputStream.read()).isEqualTo(2);
+        }
+        verify(table, never()).fileIO();
+    }
+
+    @Test
+    public void testRESTTokenFileIOSurvivesSerialization() throws Exception {
+        java.nio.file.Path sourceDirectory = 
Files.createDirectory(tempPath.resolve("rest-source"));
+        java.nio.file.Path blobFile = sourceDirectory.resolve("blob");
+        Files.write(blobFile, new byte[] {1, 2});
+
+        Identifier sourceIdentifier = 
Identifier.fromString("db.source$branch_rt");
+        String sourceRoot = "isolated://" + sourceDirectory;
+        RESTApi restApi = mock(RESTApi.class);
+        when(restApi.loadTableToken(sourceIdentifier))
+                .thenReturn(
+                        new GetTableTokenResponse(
+                                Collections.singletonMap(
+                                        IsolatedDirectoryFileIO.ROOT_DIR, 
sourceRoot),
+                                Long.MAX_VALUE));
+        RESTTokenFileIO sourceFileIO =
+                new RESTTokenFileIO(
+                        CatalogContext.create(new Options()),
+                        restApi,
+                        sourceIdentifier,
+                        new Path(sourceRoot));
+
+        FileStoreTable sourceTable = mock(FileStoreTable.class);
+        when(sourceTable.fileIO()).thenReturn(sourceFileIO);
+        Catalog catalog = mock(Catalog.class);
+        when(catalog.getTable(sourceIdentifier)).thenReturn(sourceTable);
+        CatalogLoader catalogLoader = mock(CatalogLoader.class);
+        when(catalogLoader.load()).thenReturn(catalog);
+        CatalogEnvironment catalogEnvironment = mock(CatalogEnvironment.class);
+        when(catalogEnvironment.catalogLoader()).thenReturn(catalogLoader);
+
+        FileStoreTable targetTable = mock(FileStoreTable.class);
+        when(targetTable.catalogEnvironment()).thenReturn(catalogEnvironment);
+        when(targetTable.coreOptions())
+                .thenReturn(
+                        CoreOptions.fromMap(
+                                Collections.singletonMap(
+                                        "blob-descriptor.source-table", 
"db.source$branch_rt")));
+
+        UriReaderFactory readerFactory = 
BlobDescriptorReaderFactory.create(targetTable);
+        verify(catalogLoader).load();
+        verify(catalog).getTable(sourceIdentifier);
+        // Creating the factory calls isObjectStore(), which initializes the 
dynamic token before
+        // the REST client is lost during serialization.
+        verify(restApi).loadTableToken(sourceIdentifier);
+
+        readerFactory = InstantiationUtil.clone(readerFactory);
+        String blobUri = "isolated://" + blobFile;
+        UriReaderFactory contextOnlyFactory =
+                new UriReaderFactory(CatalogContext.create(new Options()));
+        assertThatThrownBy(() -> 
contextOnlyFactory.create(blobUri).newInputStream(blobUri))
+                .isInstanceOf(UnsupportedOperationException.class)
+                .hasMessageContaining(IsolatedDirectoryFileIO.ROOT_DIR);
+        try (SeekableInputStream inputStream =
+                readerFactory.create(blobUri).newInputStream(blobUri)) {
+            assertThat(inputStream.read()).isEqualTo(1);
+            assertThat(inputStream.read()).isEqualTo(2);
+        }
+    }
+
+    @Test
+    public void testRejectSourceTableWithoutCatalogLoader() {
+        CatalogEnvironment catalogEnvironment = mock(CatalogEnvironment.class);
+        FileStoreTable targetTable = mock(FileStoreTable.class);
+        when(targetTable.catalogEnvironment()).thenReturn(catalogEnvironment);
+        when(targetTable.coreOptions())
+                .thenReturn(
+                        CoreOptions.fromMap(
+                                Collections.singletonMap(
+                                        "blob-descriptor.source-table", 
"db.source")));
+
+        assertThatThrownBy(() -> 
BlobDescriptorReaderFactory.create(targetTable))
+                .isInstanceOf(NullPointerException.class)
+                .hasMessageContaining("not supported for tables without a 
catalog loader")
+                .hasMessageContaining("external tables in REST catalogs");
+    }
+
+    @Test
+    public void testUseConfiguredExternalStorageFileIO() throws Exception {
+        java.nio.file.Path externalDirectory = 
Files.createDirectory(tempPath.resolve("external"));
+        java.nio.file.Path targetDirectory = 
Files.createDirectory(tempPath.resolve("target"));
+        java.nio.file.Path blobFile = externalDirectory.resolve("blob");
+        Files.write(blobFile, new byte[] {1, 2});
+
+        CatalogEnvironment catalogEnvironment = mock(CatalogEnvironment.class);
+        
when(catalogEnvironment.catalogContext()).thenReturn(CatalogContext.create(new 
Options()));
+        FileStoreTable targetTable = mock(FileStoreTable.class);
+        when(targetTable.fileIO()).thenReturn(isolatedFileIO(targetDirectory));
+        when(targetTable.catalogEnvironment()).thenReturn(catalogEnvironment);
+        when(targetTable.coreOptions())
+                .thenReturn(
+                        CoreOptions.fromMap(
+                                Collections.singletonMap(
+                                        "blob-descriptor.root-dir",
+                                        "isolated://" + externalDirectory)));
+
+        UriReaderFactory readerFactory = 
BlobDescriptorReaderFactory.create(targetTable);
+        String blobUri = "isolated://" + blobFile;
+        try (SeekableInputStream inputStream =
+                readerFactory.create(blobUri).newInputStream(blobUri)) {
+            assertThat(inputStream.read()).isEqualTo(1);
+            assertThat(inputStream.read()).isEqualTo(2);
+        }
+    }
+
+    private static FileIO isolatedFileIO(java.nio.file.Path root) {
+        Options options = new Options();
+        options.set(IsolatedDirectoryFileIO.ROOT_DIR, new 
Path(root.toUri()).toString());
+        IsolatedDirectoryFileIO fileIO = new IsolatedDirectoryFileIO();
+        fileIO.configure(CatalogContext.create(options));
+        return fileIO;
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkRowWrapper.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkRowWrapper.java
index 5e587f9780..079a192e7d 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkRowWrapper.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkRowWrapper.java
@@ -106,8 +106,36 @@ public class FlinkRowWrapper implements InternalRow {
             boolean checkBlobDescriptorExists,
             boolean writeNullOnFetchFailure,
             Set<Integer> blobFields) {
+        this(
+                row,
+                new UriReaderFactory(catalogContext),
+                checkBlobDescriptorExists,
+                writeNullOnFetchFailure,
+                blobFields);
+    }
+
+    public static FlinkRowWrapper fromUriReaderFactory(
+            org.apache.flink.table.data.RowData row,
+            UriReaderFactory uriReaderFactory,
+            boolean checkBlobDescriptorExists,
+            boolean writeNullOnFetchFailure,
+            Set<Integer> blobFields) {
+        return new FlinkRowWrapper(
+                row,
+                uriReaderFactory,
+                checkBlobDescriptorExists,
+                writeNullOnFetchFailure,
+                blobFields);
+    }
+
+    private FlinkRowWrapper(
+            org.apache.flink.table.data.RowData row,
+            UriReaderFactory uriReaderFactory,
+            boolean checkBlobDescriptorExists,
+            boolean writeNullOnFetchFailure,
+            Set<Integer> blobFields) {
         this.row = row;
-        this.uriReaderFactory = new UriReaderFactory(catalogContext);
+        this.uriReaderFactory = uriReaderFactory;
         this.checkBlobDescriptorExists = checkBlobDescriptorExists;
         this.writeNullOnFetchFailure = writeNullOnFetchFailure;
         this.blobFields = blobFields;
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRow.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRow.java
new file mode 100644
index 0000000000..7419e4dfbe
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRow.java
@@ -0,0 +1,217 @@
+/*
+ * 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.paimon.flink.sink;
+
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.Blob;
+import org.apache.paimon.data.BlobDescriptor;
+import org.apache.paimon.data.BlobRef;
+import org.apache.paimon.data.Decimal;
+import org.apache.paimon.data.InternalArray;
+import org.apache.paimon.data.InternalMap;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.InternalVector;
+import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.data.variant.Variant;
+import org.apache.paimon.utils.PartialRow;
+import org.apache.paimon.utils.UriReaderFactory;
+
+/** Reattaches a descriptor reader to BLOBs which passed through Flink row 
serialization. */
+final class BlobDescriptorResolvingRow extends PartialRow {
+
+    private final InternalRow wrapped;
+    private final UriReaderFactory uriReaderFactory;
+
+    BlobDescriptorResolvingRow(InternalRow wrapped, UriReaderFactory 
uriReaderFactory) {
+        super(wrapped.getFieldCount(), wrapped);
+        this.wrapped = wrapped;
+        this.uriReaderFactory = uriReaderFactory;
+    }
+
+    @Override
+    public BlobDescriptorResolvingRow replace(InternalRow row) {
+        throw new UnsupportedOperationException("Replacing the wrapped row is 
not supported.");
+    }
+
+    @Override
+    public int getFieldCount() {
+        return wrapped.getFieldCount();
+    }
+
+    @Override
+    public Blob getBlob(int pos) {
+        return withReader(super.getBlob(pos), uriReaderFactory);
+    }
+
+    private static Blob withReader(Blob blob, UriReaderFactory 
uriReaderFactory) {
+        if (!(blob instanceof BlobRef)) {
+            return blob;
+        }
+
+        BlobDescriptor descriptor = blob.toDescriptor();
+        return Blob.fromDescriptor(uriReaderFactory.create(descriptor.uri()), 
descriptor);
+    }
+
+    @Override
+    public InternalArray getArray(int pos) {
+        return new BlobDescriptorResolvingArray(super.getArray(pos), 
uriReaderFactory);
+    }
+
+    private static final class BlobDescriptorResolvingArray implements 
InternalArray {
+
+        private final InternalArray wrapped;
+        private final UriReaderFactory uriReaderFactory;
+
+        private BlobDescriptorResolvingArray(
+                InternalArray wrapped, UriReaderFactory uriReaderFactory) {
+            this.wrapped = wrapped;
+            this.uriReaderFactory = uriReaderFactory;
+        }
+
+        @Override
+        public int size() {
+            return wrapped.size();
+        }
+
+        @Override
+        public boolean isNullAt(int pos) {
+            return wrapped.isNullAt(pos);
+        }
+
+        @Override
+        public boolean getBoolean(int pos) {
+            return wrapped.getBoolean(pos);
+        }
+
+        @Override
+        public byte getByte(int pos) {
+            return wrapped.getByte(pos);
+        }
+
+        @Override
+        public short getShort(int pos) {
+            return wrapped.getShort(pos);
+        }
+
+        @Override
+        public int getInt(int pos) {
+            return wrapped.getInt(pos);
+        }
+
+        @Override
+        public long getLong(int pos) {
+            return wrapped.getLong(pos);
+        }
+
+        @Override
+        public float getFloat(int pos) {
+            return wrapped.getFloat(pos);
+        }
+
+        @Override
+        public double getDouble(int pos) {
+            return wrapped.getDouble(pos);
+        }
+
+        @Override
+        public BinaryString getString(int pos) {
+            return wrapped.getString(pos);
+        }
+
+        @Override
+        public Decimal getDecimal(int pos, int precision, int scale) {
+            return wrapped.getDecimal(pos, precision, scale);
+        }
+
+        @Override
+        public Timestamp getTimestamp(int pos, int precision) {
+            return wrapped.getTimestamp(pos, precision);
+        }
+
+        @Override
+        public byte[] getBinary(int pos) {
+            return wrapped.getBinary(pos);
+        }
+
+        @Override
+        public Variant getVariant(int pos) {
+            return wrapped.getVariant(pos);
+        }
+
+        @Override
+        public Blob getBlob(int pos) {
+            return withReader(wrapped.getBlob(pos), uriReaderFactory);
+        }
+
+        @Override
+        public InternalArray getArray(int pos) {
+            return new BlobDescriptorResolvingArray(wrapped.getArray(pos), 
uriReaderFactory);
+        }
+
+        @Override
+        public InternalVector getVector(int pos) {
+            return wrapped.getVector(pos);
+        }
+
+        @Override
+        public InternalMap getMap(int pos) {
+            return wrapped.getMap(pos);
+        }
+
+        @Override
+        public InternalRow getRow(int pos, int numFields) {
+            return wrapped.getRow(pos, numFields);
+        }
+
+        @Override
+        public boolean[] toBooleanArray() {
+            return wrapped.toBooleanArray();
+        }
+
+        @Override
+        public byte[] toByteArray() {
+            return wrapped.toByteArray();
+        }
+
+        @Override
+        public short[] toShortArray() {
+            return wrapped.toShortArray();
+        }
+
+        @Override
+        public int[] toIntArray() {
+            return wrapped.toIntArray();
+        }
+
+        @Override
+        public long[] toLongArray() {
+            return wrapped.toLongArray();
+        }
+
+        @Override
+        public float[] toFloatArray() {
+            return wrapped.toFloatArray();
+        }
+
+        @Override
+        public double[] toDoubleArray() {
+            return wrapped.toDoubleArray();
+        }
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
index 7405ae4894..b323301c79 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
@@ -29,6 +29,7 @@ import org.apache.paimon.options.MemorySize;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.table.BucketMode;
 import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.flink.api.common.RuntimeExecutionMode;
 import org.apache.flink.api.common.operators.SlotSharingGroup;
@@ -82,11 +83,17 @@ public abstract class FlinkSink<T> implements Serializable {
     protected final FileStoreTable table;
     private final boolean ignorePreviousFiles;
 
+    @Nullable private UriReaderFactory blobDescriptorReaderFactory;
+
     public FlinkSink(FileStoreTable table, boolean ignorePreviousFiles) {
         this.table = table;
         this.ignorePreviousFiles = ignorePreviousFiles;
     }
 
+    void setBlobDescriptorReaderFactory(UriReaderFactory uriReaderFactory) {
+        this.blobDescriptorReaderFactory = uriReaderFactory;
+    }
+
     public DataStreamSink<?> sinkFrom(DataStream<T> input) {
         // This commitUser is valid only for new jobs.
         // After the job starts, this commitUser will be recorded into the 
states of write and
@@ -130,18 +137,22 @@ public abstract class FlinkSink<T> implements 
Serializable {
         boolean isStreaming = isStreaming(input);
 
         boolean writeOnly = table.coreOptions().writeOnly();
+        StoreSinkWrite.Provider writeProvider =
+                StoreSinkWrite.createWriteProvider(
+                        table,
+                        env.getCheckpointConfig(),
+                        isStreaming,
+                        ignorePreviousFiles,
+                        hasSinkMaterializer(input));
+        writeProvider =
+                StoreSinkWrite.withBlobDescriptorReaderFactory(
+                        writeProvider, blobDescriptorReaderFactory);
+
         SingleOutputStreamOperator<Committable> written =
                 input.transform(
                         (writeOnly ? WRITER_WRITE_ONLY_NAME : WRITER_NAME) + " 
: " + table.name(),
                         new CommittableTypeInfo(),
-                        createWriteOperatorFactory(
-                                StoreSinkWrite.createWriteProvider(
-                                        table,
-                                        env.getCheckpointConfig(),
-                                        isStreaming,
-                                        ignorePreviousFiles,
-                                        hasSinkMaterializer(input)),
-                                commitUser));
+                        createWriteOperatorFactory(writeProvider, commitUser));
         if (parallelism == null) {
             forwardParallelism(written, input);
         } else {
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSinkBuilder.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSinkBuilder.java
index c381864747..55ed3cf3a0 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSinkBuilder.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSinkBuilder.java
@@ -33,12 +33,13 @@ import 
org.apache.paimon.flink.sink.partition.StatisticsOrRecordChannelComputer;
 import org.apache.paimon.flink.sink.partition.StatisticsOrRecordTypeInfo;
 import org.apache.paimon.flink.sorter.TableSortInfo;
 import org.apache.paimon.flink.sorter.TableSorter;
+import org.apache.paimon.table.BlobDescriptorReaderFactory;
 import org.apache.paimon.table.BucketMode;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.PostponeUtils;
 import org.apache.paimon.table.Table;
 import org.apache.paimon.table.sink.ChannelComputer;
-import org.apache.paimon.utils.BlobDescriptorUtils;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.flink.api.common.functions.MapFunction;
 import org.apache.flink.streaming.api.datastream.DataStream;
@@ -72,6 +73,7 @@ import static 
org.apache.paimon.flink.utils.ParallelismUtils.forwardParallelism;
 import static org.apache.paimon.flink.utils.ParallelismUtils.setParallelism;
 import static org.apache.paimon.table.BucketMode.BUCKET_UNAWARE;
 import static org.apache.paimon.utils.Preconditions.checkArgument;
+import static org.apache.paimon.utils.Preconditions.checkNotNull;
 import static org.apache.paimon.utils.Preconditions.checkState;
 
 /**
@@ -90,6 +92,7 @@ public class FlinkSinkBuilder {
     @Nullable protected Map<String, String> overwritePartition;
     @Nullable private Integer parallelism;
     @Nullable private TableSortInfo tableSortInfo;
+    @Nullable private UriReaderFactory blobDescriptorReaderFactory;
 
     // ============== for extension ==============
 
@@ -213,16 +216,14 @@ public class FlinkSinkBuilder {
     public DataStreamSink<?> build() {
         setParallelismIfAdaptiveConflict();
         input = trySortInput(input);
-        CatalogContext contextForDescriptor =
-                BlobDescriptorUtils.getCatalogContext(
-                        table.catalogEnvironment().catalogContext(),
-                        table.coreOptions().toConfiguration());
+        UriReaderFactory readerFactoryForDescriptor = 
BlobDescriptorReaderFactory.create(table);
+        blobDescriptorReaderFactory = readerFactoryForDescriptor;
 
         DataStream<InternalRow> input =
-                mapToInternalRow(
+                mapToInternalRowWithUriReaderFactory(
                         this.input,
                         table.rowType(),
-                        contextForDescriptor,
+                        readerFactoryForDescriptor,
                         table.coreOptions().blobWriteNullOnMissingFile(),
                         table.coreOptions().blobWriteNullOnFetchFailure());
         if (table.coreOptions().localMergeEnabled() && 
table.schema().primaryKeys().size() > 0) {
@@ -274,6 +275,20 @@ public class FlinkSinkBuilder {
             CatalogContext catalogContext,
             boolean checkBlobDescriptorExists,
             boolean writeNullOnFetchFailure) {
+        return mapToInternalRowWithUriReaderFactory(
+                input,
+                rowType,
+                new UriReaderFactory(catalogContext),
+                checkBlobDescriptorExists,
+                writeNullOnFetchFailure);
+    }
+
+    private static DataStream<InternalRow> 
mapToInternalRowWithUriReaderFactory(
+            DataStream<RowData> input,
+            org.apache.paimon.types.RowType rowType,
+            UriReaderFactory uriReaderFactory,
+            boolean checkBlobDescriptorExists,
+            boolean writeNullOnFetchFailure) {
         Set<Integer> blobFields =
                 checkBlobDescriptorExists
                         ? FlinkRowWrapper.blobFieldIndexes(rowType)
@@ -282,9 +297,9 @@ public class FlinkSinkBuilder {
                 input.map(
                                 (MapFunction<RowData, InternalRow>)
                                         r ->
-                                                new FlinkRowWrapper(
+                                                
FlinkRowWrapper.fromUriReaderFactory(
                                                         r,
-                                                        catalogContext,
+                                                        uriReaderFactory,
                                                         
checkBlobDescriptorExists,
                                                         
writeNullOnFetchFailure,
                                                         blobFields))
@@ -297,14 +312,22 @@ public class FlinkSinkBuilder {
 
     protected DataStreamSink<?> buildDynamicBucketSink(
             DataStream<InternalRow> input, boolean globalIndex) {
-        return compactSink && !globalIndex
-                // todo support global index sort compact
-                ? new DynamicBucketCompactSink(table, 
overwritePartition).build(input, parallelism)
-                : globalIndex
-                        ? new GlobalDynamicBucketSink(table, 
overwritePartition)
-                                .build(input, parallelism)
-                        : new RowDynamicBucketSink(table, overwritePartition)
-                                .build(input, parallelism);
+        if (compactSink && !globalIndex) {
+            // todo support global index sort compact
+            return configureBlobDescriptorReaderFactory(
+                            new DynamicBucketCompactSink(table, 
overwritePartition))
+                    .build(input, parallelism);
+        }
+
+        if (globalIndex) {
+            return configureBlobDescriptorReaderFactory(
+                            new GlobalDynamicBucketSink(table, 
overwritePartition))
+                    .build(input, parallelism);
+        }
+
+        return configureBlobDescriptorReaderFactory(
+                        new RowDynamicBucketSink(table, overwritePartition))
+                .build(input, parallelism);
     }
 
     protected DataStreamSink<?> buildForFixedBucket(DataStream<InternalRow> 
input) {
@@ -320,8 +343,8 @@ public class FlinkSinkBuilder {
         }
         DataStream<InternalRow> partitioned =
                 partition(input, new RowDataChannelComputer(table.schema()), 
parallelism);
-        FixedBucketSink sink = new FixedBucketSink(table, overwritePartition);
-        return sink.sinkFrom(partitioned);
+        return configureBlobDescriptorReaderFactory(new FixedBucketSink(table, 
overwritePartition))
+                .sinkFrom(partitioned);
     }
 
     private DataStreamSink<?> buildPostponeBucketSink(DataStream<InternalRow> 
input) {
@@ -334,8 +357,9 @@ public class FlinkSinkBuilder {
                 channelComputer = new 
PostponeBucketChannelComputer(table.schema());
             }
             DataStream<InternalRow> partitioned = partition(input, 
channelComputer, parallelism);
-            PostponeBucketSink sink = new PostponeBucketSink(table, 
overwritePartition);
-            return sink.sinkFrom(partitioned);
+            return configureBlobDescriptorReaderFactory(
+                            new PostponeBucketSink(table, overwritePartition))
+                    .sinkFrom(partitioned);
         } else {
             Map<BinaryRow, Integer> knownNumBuckets = 
PostponeUtils.getKnownNumBuckets(table);
             DataStream<InternalRow> partitioned =
@@ -346,9 +370,10 @@ public class FlinkSinkBuilder {
 
             FileStoreTable tableForWrite = 
PostponeUtils.tableForFixBucketWrite(table);
 
-            PostponeFixedBucketSink sink =
-                    new PostponeFixedBucketSink(tableForWrite, 
overwritePartition, knownNumBuckets);
-            return sink.sinkFrom(partitioned);
+            return configureBlobDescriptorReaderFactory(
+                            new PostponeFixedBucketSink(
+                                    tableForWrite, overwritePartition, 
knownNumBuckets))
+                    .sinkFrom(partitioned);
         }
     }
 
@@ -370,7 +395,15 @@ public class FlinkSinkBuilder {
             }
         }
 
-        return new RowAppendTableSink(table, overwritePartition, 
parallelism).sinkFrom(input);
+        return configureBlobDescriptorReaderFactory(
+                        new RowAppendTableSink(table, overwritePartition, 
parallelism))
+                .sinkFrom(input);
+    }
+
+    private <T extends FlinkSink<?>> T configureBlobDescriptorReaderFactory(T 
sink) {
+        sink.setBlobDescriptorReaderFactory(
+                checkNotNull(blobDescriptorReaderFactory, "BLOB descriptor 
reader is not set."));
+        return sink;
     }
 
     private DataStream<InternalRow> 
applyDynamicPartitionShuffle(DataStream<InternalRow> input) {
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWrite.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWrite.java
index 06c0353f18..5aee320f6e 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWrite.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWrite.java
@@ -30,6 +30,7 @@ import org.apache.paimon.table.sink.SinkRecord;
 import org.apache.paimon.table.sink.TableWriteImpl;
 import org.apache.paimon.utils.Preconditions;
 import org.apache.paimon.utils.SerializableRunnable;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.runtime.io.disk.iomanager.IOManager;
@@ -50,6 +51,8 @@ public interface StoreSinkWrite {
 
     void setWriteRestore(WriteRestore writeRestore);
 
+    default void setBlobDescriptorReaderFactory(UriReaderFactory 
uriReaderFactory) {}
+
     @Nullable
     SinkRecord write(InternalRow rowData) throws Exception;
 
@@ -97,6 +100,21 @@ public interface StoreSinkWrite {
                 @Nullable MetricGroup metricGroup);
     }
 
+    static Provider withBlobDescriptorReaderFactory(
+            Provider provider, @Nullable UriReaderFactory uriReaderFactory) {
+        if (uriReaderFactory == null) {
+            return provider;
+        }
+
+        return (table, commitUser, state, ioManager, memoryPoolFactory, 
metricGroup) -> {
+            StoreSinkWrite write =
+                    provider.provide(
+                            table, commitUser, state, ioManager, 
memoryPoolFactory, metricGroup);
+            write.setBlobDescriptorReaderFactory(uriReaderFactory);
+            return write;
+        };
+    }
+
     static StoreSinkWrite.Provider createWriteProvider(
             FileStoreTable fileStoreTable,
             CheckpointConfig checkpointConfig,
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWriteImpl.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWriteImpl.java
index 41c23f5002..e314b997c5 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWriteImpl.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreSinkWriteImpl.java
@@ -30,6 +30,7 @@ import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.table.sink.SinkRecord;
 import org.apache.paimon.table.sink.TableWriteImpl;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.runtime.io.disk.iomanager.IOManager;
@@ -57,6 +58,8 @@ public class StoreSinkWriteImpl implements StoreSinkWrite {
     private final MemoryPoolFactory memoryPoolFactory;
     @Nullable private final MetricGroup metricGroup;
 
+    @Nullable private UriReaderFactory blobDescriptorReaderFactory;
+
     protected TableWriteImpl<?> write;
 
     public StoreSinkWriteImpl(
@@ -102,16 +105,27 @@ public class StoreSinkWriteImpl implements StoreSinkWrite 
{
         write.withWriteRestore(writeRestore);
     }
 
+    @Override
+    public void setBlobDescriptorReaderFactory(UriReaderFactory 
uriReaderFactory) {
+        this.blobDescriptorReaderFactory = uriReaderFactory;
+    }
+
     @Override
     @Nullable
     public SinkRecord write(InternalRow rowData) throws Exception {
-        return write.writeAndReturn(rowData);
+        return write.writeAndReturn(withBlobDescriptorReader(rowData));
     }
 
     @Override
     @Nullable
     public SinkRecord write(InternalRow rowData, int bucket) throws Exception {
-        return write.writeAndReturn(rowData, bucket);
+        return write.writeAndReturn(withBlobDescriptorReader(rowData), bucket);
+    }
+
+    private InternalRow withBlobDescriptorReader(InternalRow rowData) {
+        return blobDescriptorReaderFactory == null
+                ? rowData
+                : new BlobDescriptorResolvingRow(rowData, 
blobDescriptorReaderFactory);
     }
 
     @Override
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BlobTableITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BlobTableITCase.java
index e2ad151fdc..d85c534f79 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BlobTableITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BlobTableITCase.java
@@ -386,6 +386,36 @@ public class BlobTableITCase extends CatalogITCaseBase {
                 .containsExactlyInAnyOrder(Row.of(1, "paimon", blobData));
     }
 
+    @Test
+    public void testMaterializeDescriptorWithSourceTableFileIO() {
+        tEnv.executeSql(
+                "CREATE TABLE blob_descriptor_source (id INT, picture BYTES)"
+                        + " WITH ('row-tracking.enabled'='true',"
+                        + " 'data-evolution.enabled'='true',"
+                        + " 'blob-field'='picture')");
+        batchSql("INSERT INTO blob_descriptor_source VALUES" + " (1, 
X'48656C6C6F'), (2, X'5945')");
+        batchSql("ALTER TABLE blob_descriptor_source SET 
('blob-as-descriptor'='true')");
+
+        String sourceTable = tEnv.getCurrentDatabase() + 
".blob_descriptor_source";
+        tEnv.executeSql(
+                "CREATE TABLE blob_descriptor_target ("
+                        + "id INT, picture BYTES, PRIMARY KEY (id) NOT 
ENFORCED)"
+                        + " WITH ('bucket'='2',"
+                        + " 'blob-field'='picture',"
+                        + " 'blob-descriptor.source-table'='"
+                        + sourceTable
+                        + "')");
+        batchSql(
+                "INSERT INTO blob_descriptor_target "
+                        + "/*+ OPTIONS('sink.parallelism' = '2') */ "
+                        + "SELECT * FROM blob_descriptor_source");
+
+        assertThat(batchSql("SELECT * FROM blob_descriptor_target ORDER BY 
id"))
+                .containsExactly(
+                        Row.of(1, new byte[] {72, 101, 108, 108, 111}),
+                        Row.of(2, new byte[] {89, 69}));
+    }
+
     @Test
     public void testWriteBlobWithBuiltInFunction() throws Exception {
         byte[] blobData = new byte[1024 * 1024];
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkRowWrapperTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkRowWrapperTest.java
index 3141ef7062..060c90d309 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkRowWrapperTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/FlinkRowWrapperTest.java
@@ -21,7 +21,9 @@ package org.apache.paimon.flink;
 import org.apache.paimon.catalog.CatalogContext;
 import org.apache.paimon.data.Blob;
 import org.apache.paimon.data.BlobDescriptor;
+import org.apache.paimon.fs.IsolatedDirectoryFileIO;
 import org.apache.paimon.options.Options;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import com.sun.net.httpserver.HttpExchange;
 import com.sun.net.httpserver.HttpServer;
@@ -86,6 +88,30 @@ public class FlinkRowWrapperTest {
         assertThat(wrapper.getBlob(0).toData()).isEqualTo(bytes);
     }
 
+    @Test
+    public void testReadBlobWithProvidedUriReaderFactory() throws Exception {
+        byte[] bytes = new byte[] {1, 2, 3};
+        java.nio.file.Path blobFile = 
tempPath.resolve("provided-file-io.blob");
+        Files.write(blobFile, bytes);
+        String blobUri = "isolated://" + blobFile;
+
+        Options options = new Options();
+        options.set(IsolatedDirectoryFileIO.ROOT_DIR, "isolated://" + 
tempPath);
+        IsolatedDirectoryFileIO fileIO = new IsolatedDirectoryFileIO();
+        fileIO.configure(CatalogContext.create(options));
+        UriReaderFactory readerFactory = UriReaderFactory.fromFileIO(fileIO);
+
+        FlinkRowWrapper wrapper =
+                FlinkRowWrapper.fromUriReaderFactory(
+                        descriptorRow(blobUri, bytes.length),
+                        readerFactory,
+                        false,
+                        false,
+                        Collections.singleton(0));
+
+        assertThat(wrapper.getBlob(0).toData()).isEqualTo(bytes);
+    }
+
     @Test
     public void testMissingHttpBlobDescriptorWithNonBlobColumnBefore() throws 
Exception {
         httpServer.createContext(
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRowTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRowTest.java
new file mode 100644
index 0000000000..d80389276e
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BlobDescriptorResolvingRowTest.java
@@ -0,0 +1,74 @@
+/*
+ * 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.paimon.flink.sink;
+
+import org.apache.paimon.data.Blob;
+import org.apache.paimon.data.GenericArray;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.flink.utils.InternalRowTypeSerializer;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.UriReaderFactory;
+
+import org.apache.flink.core.memory.DataInputDeserializer;
+import org.apache.flink.core.memory.DataOutputSerializer;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Files;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link BlobDescriptorResolvingRow}. */
+class BlobDescriptorResolvingRowTest {
+
+    @TempDir java.nio.file.Path tempPath;
+
+    @Test
+    void testArrayBlobAfterFlinkSerialization() throws Exception {
+        byte[] expected = new byte[] {1, 2, 3};
+        java.nio.file.Path blobPath = tempPath.resolve("blob");
+        Files.write(blobPath, expected);
+
+        RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
+        GenericRow row =
+                GenericRow.of(
+                        new GenericArray(
+                                new Object[] {
+                                    Blob.fromFile(LocalFileIO.create(), 
blobPath.toUri().toString())
+                                }));
+        InternalRow serialized = serializeAndDeserialize(row, rowType);
+
+        BlobDescriptorResolvingRow resolvingRow =
+                new BlobDescriptorResolvingRow(
+                        serialized, 
UriReaderFactory.fromFileIO(LocalFileIO.create()));
+
+        
assertThat(resolvingRow.getArray(0).getBlob(0).toData()).isEqualTo(expected);
+    }
+
+    private static InternalRow serializeAndDeserialize(InternalRow row, 
RowType rowType)
+            throws Exception {
+        InternalRowTypeSerializer serializer = new 
InternalRowTypeSerializer(rowType);
+        DataOutputSerializer output = new DataOutputSerializer(100);
+        serializer.serialize(row, output);
+        return serializer.deserialize(new 
DataInputDeserializer(output.wrapAsByteBuffer()));
+    }
+}
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRowWrapper.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRowWrapper.java
index f6bc91170b..7823839816 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRowWrapper.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRowWrapper.java
@@ -65,7 +65,7 @@ public class SparkInternalRowWrapper implements InternalRow, 
Serializable {
     private transient org.apache.spark.sql.catalyst.InternalRow internalRow;
 
     public SparkInternalRowWrapper(StructType tableSchema, int length) {
-        this(tableSchema, length, null, null);
+        this(tableSchema, length, null, (CatalogContext) null);
     }
 
     public SparkInternalRowWrapper(
@@ -73,12 +73,28 @@ public class SparkInternalRowWrapper implements 
InternalRow, Serializable {
             int length,
             StructType dataSchema,
             CatalogContext catalogContext) {
+        this(tableSchema, length, dataSchema, new 
UriReaderFactory(catalogContext));
+    }
+
+    public static SparkInternalRowWrapper fromUriReaderFactory(
+            StructType tableSchema,
+            int length,
+            StructType dataSchema,
+            @Nullable UriReaderFactory uriReaderFactory) {
+        return new SparkInternalRowWrapper(tableSchema, length, dataSchema, 
uriReaderFactory);
+    }
+
+    private SparkInternalRowWrapper(
+            StructType tableSchema,
+            int length,
+            StructType dataSchema,
+            @Nullable UriReaderFactory uriReaderFactory) {
         this.tableSchema = tableSchema;
         this.length = length;
         this.dataSchema = dataSchema;
         this.fieldIndexMap =
                 dataSchema != null ? buildFieldIndexMap(tableSchema, 
dataSchema) : null;
-        this.uriReaderFactory = new UriReaderFactory(catalogContext);
+        this.uriReaderFactory = uriReaderFactory;
     }
 
     public SparkInternalRowWrapper 
replace(org.apache.spark.sql.catalyst.InternalRow internalRow) {
@@ -305,10 +321,11 @@ public class SparkInternalRowWrapper implements 
InternalRow, Serializable {
         if (dataSchema != null) {
             StructType nestedDataSchema = (StructType) 
dataSchema.fields()[actualPos].dataType();
             int dataNumFields = nestedDataSchema.size();
-            return new SparkInternalRowWrapper(nestedTableSchema, numFields, 
nestedDataSchema, null)
+            return new SparkInternalRowWrapper(
+                            nestedTableSchema, numFields, nestedDataSchema, 
uriReaderFactory)
                     .replace(internalRow.getStruct(actualPos, dataNumFields));
         }
-        return new SparkInternalRowWrapper(nestedTableSchema, numFields)
+        return new SparkInternalRowWrapper(nestedTableSchema, numFields, null, 
uriReaderFactory)
                 .replace(internalRow.getStruct(actualPos, numFields));
     }
 
@@ -483,7 +500,8 @@ public class SparkInternalRowWrapper implements 
InternalRow, Serializable {
 
         @Override
         public InternalRow getRow(int pos, int numFields) {
-            return new SparkInternalRowWrapper((StructType) elementType, 
numFields)
+            return new SparkInternalRowWrapper(
+                            (StructType) elementType, numFields, null, 
uriReaderFactory)
                     .replace(arrayData.getStruct(pos, numFields));
         }
     }
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkRow.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkRow.java
index dee24ab82d..643fc016d9 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkRow.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkRow.java
@@ -65,14 +65,23 @@ public class SparkRow implements InternalRow, Serializable {
     private final UriReaderFactory uriReaderFactory;
 
     public SparkRow(RowType type, Row row) {
-        this(type, row, RowKind.INSERT, null);
+        this(type, row, RowKind.INSERT, (CatalogContext) null);
     }
 
     public SparkRow(RowType type, Row row, RowKind rowkind, CatalogContext 
catalogContext) {
+        this(type, row, rowkind, new UriReaderFactory(catalogContext));
+    }
+
+    public static SparkRow fromUriReaderFactory(
+            RowType type, Row row, RowKind rowkind, UriReaderFactory 
uriReaderFactory) {
+        return new SparkRow(type, row, rowkind, uriReaderFactory);
+    }
+
+    private SparkRow(RowType type, Row row, RowKind rowkind, UriReaderFactory 
uriReaderFactory) {
         this.type = type;
         this.row = row;
         this.rowKind = rowkind;
-        this.uriReaderFactory = new UriReaderFactory(catalogContext);
+        this.uriReaderFactory = uriReaderFactory;
     }
 
     @Override
@@ -185,7 +194,8 @@ public class SparkRow implements InternalRow, Serializable {
 
     @Override
     public InternalRow getRow(int i, int i1) {
-        return new SparkRow((RowType) type.getTypeAt(i), row.getStruct(i));
+        return new SparkRow(
+                (RowType) type.getTypeAt(i), row.getStruct(i), RowKind.INSERT, 
uriReaderFactory);
     }
 
     private static int toPaimonDate(Object object) {
@@ -376,7 +386,7 @@ public class SparkRow implements InternalRow, Serializable {
 
         @Override
         public InternalRow getRow(int i, int i1) {
-            return new SparkRow((RowType) elementType, getAs(i));
+            return new SparkRow((RowType) elementType, getAs(i), 
RowKind.INSERT, uriReaderFactory);
         }
 
         @Override
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/write/PaimonV2MetadataAwareDataWriter.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/write/PaimonV2MetadataAwareDataWriter.java
index fd6786a885..7635756986 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/write/PaimonV2MetadataAwareDataWriter.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/write/PaimonV2MetadataAwareDataWriter.java
@@ -19,9 +19,9 @@
 package org.apache.paimon.spark.write;
 
 import org.apache.paimon.CoreOptions;
-import org.apache.paimon.catalog.CatalogContext;
 import org.apache.paimon.table.sink.BatchWriteBuilder;
 import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.spark.sql.catalyst.InternalRow;
 import org.apache.spark.sql.types.StructType;
@@ -42,14 +42,14 @@ public class PaimonV2MetadataAwareDataWriter extends 
PaimonV2DataWriter {
             StructType dataSchema,
             StructType metadataSchema,
             CoreOptions coreOptions,
-            CatalogContext catalogContext,
+            UriReaderFactory uriReaderFactory,
             RowType paimonWriteType) {
         super(
                 writeBuilder,
                 rowTrackingWriteSchema,
                 dataSchema,
                 coreOptions,
-                catalogContext,
+                uriReaderFactory,
                 Option.empty(),
                 Option.apply(paimonWriteType),
                 Option.apply(metadataSchema),
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionPaimonWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionPaimonWriter.scala
index 12f7159161..e99afc95b6 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionPaimonWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionPaimonWriter.scala
@@ -48,6 +48,7 @@ case class DataEvolutionPaimonWriter(paimonTable: 
FileStoreTable, dataSplits: Se
       columnNames: Seq[String],
       rawBlobPlaceholderMarkerColumns: Map[String, String] = Map.empty): 
Seq[CommitMessage] = {
     val sparkSession = data.sparkSession
+    val uriReaderFactory = uriReaderFactoryForBlobDescriptor
     import sparkSession.implicits._
     assert(data.columns.length == columnNames.size + 2 + 
rawBlobPlaceholderMarkerColumns.size)
     val writeType = table.rowType().project(columnNames.asJava)
@@ -110,7 +111,7 @@ case class DataEvolutionPaimonWriter(paimonTable: 
FileStoreTable, dataSplits: Se
               writeBuilder,
               writeType,
               firstRowIdToPartitionMapBroadcast.value,
-              catalogContextForBlobDescriptor,
+              uriReaderFactory,
               rawBlobPlaceholderMarkerIndexes)
             try {
               iter.foreach(row => write.write(row))
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
index 92df5a7f06..1518a77934 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/PaimonSparkWriter.scala
@@ -41,7 +41,7 @@ import org.apache.paimon.table.{FileStoreTable, 
PostponeUtils, SpecialFields}
 import org.apache.paimon.table.BucketMode._
 import org.apache.paimon.table.sink._
 import org.apache.paimon.types.{RowKind, RowType}
-import org.apache.paimon.utils.SerializationUtils
+import org.apache.paimon.utils.{SerializationUtils, UriReaderFactory}
 
 import org.apache.spark.{Partitioner, TaskContext}
 import org.apache.spark.rdd.RDD
@@ -107,6 +107,7 @@ case class PaimonSparkWriter(
 
   def write(data: DataFrame): Seq[CommitMessage] = {
     val sparkSession = data.sparkSession
+    val uriReaderFactory = uriReaderFactoryForBlobDescriptor
     import sparkSession.implicits._
 
     val withInitBucketCol = bucketMode match {
@@ -138,7 +139,7 @@ case class PaimonSparkWriter(
       writeRowTracking,
       fullCompactionDeltaCommits,
       batchId,
-      catalogContextForBlobDescriptor,
+      uriReaderFactory,
       postponePartitionBucketComputer
     )
 
@@ -210,7 +211,11 @@ case class PaimonSparkWriter(
           .map(_.toInt)
           .getOrElse(sparkParallelism)
         val bootstrapped =
-          bootstrapAndRepartitionByKeyHash(withInitBucketCol, 
assignerParallelism, rowKindColIdx)
+          bootstrapAndRepartitionByKeyHash(
+            withInitBucketCol,
+            assignerParallelism,
+            rowKindColIdx,
+            uriReaderFactory)
 
         val globalDynamicBucketProcessor =
           GlobalDynamicBucketProcessor(
@@ -245,7 +250,8 @@ case class PaimonSparkWriter(
             sparkSession,
             withInitBucketCol,
             assignerParallelism,
-            numAssigners)
+            numAssigners,
+            uriReaderFactory)
         }
 
         if (table.snapshotManager().latestSnapshotFromFileSystem() == null) {
@@ -264,7 +270,7 @@ case class PaimonSparkWriter(
                 )
               row => {
                 val sparkRow =
-                  new SparkRow(writeType, row, RowKind.INSERT, 
catalogContextForBlobDescriptor)
+                  SparkRow.fromUriReaderFactory(writeType, row, 
RowKind.INSERT, uriReaderFactory)
                 assigner.assign(
                   extractor.partition(sparkRow),
                   extractor.trimmedPrimaryKey(sparkRow).hashCode)
@@ -445,7 +451,8 @@ case class PaimonSparkWriter(
   private def bootstrapAndRepartitionByKeyHash(
       data: DataFrame,
       parallelism: Int,
-      rowKindColIdx: Int): RDD[(KeyPartOrRow, Array[Byte])] = {
+      rowKindColIdx: Int,
+      uriReaderFactory: UriReaderFactory): RDD[(KeyPartOrRow, Array[Byte])] = {
     val numSparkPartitions = data.rdd.getNumPartitions
     val primaryKeys = table.schema().primaryKeys()
     val bootstrapType = IndexBootstrap.bootstrapType(table.schema())
@@ -464,7 +471,7 @@ case class PaimonSparkWriter(
               .toCloseableIterator
             TaskContext.get().addTaskCompletionListener[Unit](_ => 
bootstrapIterator.close())
             val toPaimonRow =
-              SparkRowUtils.toPaimonRow(rowType, rowKindColIdx, 
catalogContextForBlobDescriptor)
+              SparkRowUtils.toPaimonRow(rowType, rowKindColIdx, 
uriReaderFactory)
 
             bootstrapIterator.asScala
               .map(
@@ -490,7 +497,8 @@ case class PaimonSparkWriter(
       sparkSession: SparkSession,
       data: DataFrame,
       parallelism: Int,
-      numAssigners: Int): DataFrame = {
+      numAssigners: Int,
+      uriReaderFactory: UriReaderFactory): DataFrame = {
     sparkSession.createDataFrame(
       data.rdd
         .mapPartitions(
@@ -499,7 +507,7 @@ case class PaimonSparkWriter(
             iterator.map(
               row => {
                 val sparkRow =
-                  new SparkRow(writeType, row, RowKind.INSERT, 
catalogContextForBlobDescriptor)
+                  SparkRow.fromUriReaderFactory(writeType, row, 
RowKind.INSERT, uriReaderFactory)
                 val partitionHash = 
rowPartitionKeyExtractor.partition(sparkRow).hashCode
                 val keyHash = 
rowPartitionKeyExtractor.trimmedPrimaryKey(sparkRow).hashCode
                 (
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkPostponeCompactProcedure.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkPostponeCompactProcedure.scala
index 199702e088..0e136e0e97 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkPostponeCompactProcedure.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkPostponeCompactProcedure.scala
@@ -29,9 +29,9 @@ import org.apache.paimon.spark.commands.{EncoderSerDeGroup, 
PostponeFixBucketPro
 import org.apache.paimon.spark.schema.SparkSystemColumns.{BUCKET_COL, 
ROW_KIND_COL}
 import org.apache.paimon.spark.util.{ScanPlanHelper, SparkRowUtils}
 import org.apache.paimon.spark.write.{PaimonDataWrite, WriteTaskResult}
-import org.apache.paimon.table.{BucketMode, FileStoreTable, PostponeUtils}
+import org.apache.paimon.table.{BlobDescriptorReaderFactory, BucketMode, 
FileStoreTable, PostponeUtils}
 import org.apache.paimon.table.sink.{CommitMessage, CommitMessageImpl}
-import org.apache.paimon.utils.{BlobDescriptorUtils, SerializationUtils}
+import org.apache.paimon.utils.{SerializationUtils, UriReaderFactory}
 
 import org.apache.spark.HashPartitioner
 import org.apache.spark.rdd.RDD
@@ -85,14 +85,10 @@ case class SparkPostponeCompactProcedure(
   private def newDataWrite(
       realTable: FileStoreTable,
       rowKindColIdx: Int,
-      postponePartitionBucketComputer: 
SparkPostponeCompactProcedure.PostponePartitionBucketComputer)
-      : PaimonDataWrite = {
+      postponePartitionBucketComputer: 
SparkPostponeCompactProcedure.PostponePartitionBucketComputer,
+      uriReaderFactoryForBlobDescriptor: UriReaderFactory): PaimonDataWrite = {
     val rowType = table.rowType()
     val coreOptions = table.coreOptions()
-    val catalogContextForBlobDescriptor =
-      BlobDescriptorUtils.getCatalogContext(
-        table.catalogEnvironment().catalogContext(),
-        coreOptions.toConfiguration)
 
     val dataWrite = PaimonDataWrite(
       realTable.newBatchWriteBuilder,
@@ -101,7 +97,7 @@ case class SparkPostponeCompactProcedure(
       writeRowTracking = coreOptions.dataEvolutionEnabled(),
       Option.apply(coreOptions.fullCompactionDeltaCommits()),
       None,
-      catalogContextForBlobDescriptor,
+      uriReaderFactoryForBlobDescriptor,
       Some(postponePartitionBucketComputer)
     )
     dataWrite
@@ -146,6 +142,7 @@ case class SparkPostponeCompactProcedure(
       LOG.info("Postpone bucket and real Level-0 buckets are empty, no compact 
job to execute.")
       return
     }
+    val uriReaderFactory = BlobDescriptorReaderFactory.create(table)
 
     val rowWorkAndKind: 
(RDD[SparkPostponeCompactProcedure.PostponeCompactWork], Int) =
       if (splits.isEmpty) {
@@ -175,13 +172,10 @@ case class SparkPostponeCompactProcedure(
           .toDF()
         val rowKindColIdx = 
SparkRowUtils.getFieldIndex(withInitBucketCol.schema, ROW_KIND_COL)
         val rowType = table.rowType()
-        val catalogContext = BlobDescriptorUtils.getCatalogContext(
-          table.catalogEnvironment().catalogContext(),
-          table.coreOptions().toConfiguration)
         val rowWorks = dataFrame.rdd.mapPartitions {
           rows =>
             val extractor = realTable.createRowKeyExtractor()
-            val toPaimonRow = SparkRowUtils.toPaimonRow(rowType, 
rowKindColIdx, catalogContext)
+            val toPaimonRow = SparkRowUtils.toPaimonRow(rowType, 
rowKindColIdx, uriReaderFactory)
             rows.map {
               row =>
                 extractor.setRecord(toPaimonRow(row))
@@ -229,7 +223,11 @@ case class SparkPostponeCompactProcedure(
           Iterator.empty
         } else {
           val dataWrite =
-            newDataWrite(realTable, rowWorkAndKind._2, 
postponePartitionBucketComputer)
+            newDataWrite(
+              realTable,
+              rowWorkAndKind._2,
+              postponePartitionBucketComputer,
+              uriReaderFactory)
           dataWrite.write.withWriteRestore(
             new FileSystemWriteRestore(
               realTable.coreOptions(),
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkRowUtils.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkRowUtils.scala
index eee5545454..c50db667b6 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkRowUtils.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/util/SparkRowUtils.scala
@@ -21,6 +21,7 @@ package org.apache.paimon.spark.util
 import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.spark.SparkRow
 import org.apache.paimon.types.{RowKind, RowType}
+import org.apache.paimon.utils.UriReaderFactory
 
 import org.apache.spark.sql.Row
 import org.apache.spark.sql.types.StructType
@@ -31,14 +32,23 @@ object SparkRowUtils {
       writeType: RowType,
       rowkindColIdx: Int,
       catalogContext: CatalogContext): Row => SparkRow = {
+    toPaimonRow(writeType, rowkindColIdx, new UriReaderFactory(catalogContext))
+  }
+
+  def toPaimonRow(
+      writeType: RowType,
+      rowkindColIdx: Int,
+      uriReaderFactory: UriReaderFactory): Row => SparkRow = {
     if (rowkindColIdx != -1) {
       row =>
-        new SparkRow(
+        SparkRow.fromUriReaderFactory(
           writeType,
           row,
           RowKind.fromByteValue(row.getByte(rowkindColIdx)),
-          catalogContext)
-    } else { row => new SparkRow(writeType, row, RowKind.INSERT, 
catalogContext) }
+          uriReaderFactory)
+    } else {
+      row => SparkRow.fromUriReaderFactory(writeType, row, RowKind.INSERT, 
uriReaderFactory)
+    }
   }
 
   def getFieldIndex(schema: StructType, colName: String): Int = {
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
index 67d5c93ab9..c56e742d5e 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/DataEvolutionTableDataWrite.scala
@@ -19,7 +19,6 @@
 package org.apache.paimon.spark.write
 
 import org.apache.paimon.casting.FallbackMappingRow
-import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.data.{BinaryRow, BlobArrayPlaceholder, 
BlobPlaceholder, GenericRow, InternalRow}
 import org.apache.paimon.data.serializer.InternalSerializers
 import org.apache.paimon.disk.IOManager
@@ -33,6 +32,7 @@ import org.apache.paimon.types.{DataTypeRoot, RowType}
 import org.apache.paimon.types.VectorType.isVectorStoreFile
 import org.apache.paimon.utils.RecordWriter
 import org.apache.paimon.utils.SerializationUtils
+import org.apache.paimon.utils.UriReaderFactory
 
 import org.apache.spark.sql.Row
 import org.slf4j.LoggerFactory
@@ -47,7 +47,7 @@ case class DataEvolutionTableDataWrite(
     writeBuilder: BatchWriteBuilder,
     writeType: RowType,
     firstRowIdToPartitionMap: mutable.HashMap[Long, (Array[Byte], Long)],
-    catalogContext: CatalogContext,
+    uriReaderFactory: UriReaderFactory,
     rawBlobPlaceholderMarkerIndexes: Map[Int, Int])
   extends InnerTableV1DataWrite {
 
@@ -58,7 +58,7 @@ case class DataEvolutionTableDataWrite(
   private val commitMessages = ListBuffer[CommitMessageImpl]()
 
   private val toPaimonRow = {
-    SparkRowUtils.toPaimonRow(writeType, -1, catalogContext)
+    SparkRowUtils.toPaimonRow(writeType, -1, uriReaderFactory)
   }
   private lazy val rowSerializer = InternalSerializers.create(writeType)
   private val rawBlobFallbackFields = 
rawBlobPlaceholderMarkerIndexes.toSeq.sortBy(_._1).toArray
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonBatchWriteBase.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonBatchWriteBase.scala
index 5ac1688223..302ae397d5 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonBatchWriteBase.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonBatchWriteBase.scala
@@ -27,6 +27,7 @@ import org.apache.paimon.spark.rowops.PaimonCopyOnWriteScan
 import org.apache.paimon.spark.schema.PaimonMetadataColumn.{FILE_PATH, ROW_ID, 
SEQUENCE_NUMBER}
 import org.apache.paimon.table.{FileStoreTable, SpecialFields}
 import org.apache.paimon.table.sink.{BatchWriteBuilder, CommitMessage, 
CommitMessageImpl}
+import org.apache.paimon.utils.UriReaderFactory
 
 import org.apache.spark.sql.connector.metric.CustomTaskMetric
 import org.apache.spark.sql.connector.write.{DataWriterFactory, 
PhysicalWriteInfo, WriterCommitMessage}
@@ -81,22 +82,23 @@ abstract class PaimonBatchWriteBase(
     StructType(Seq(FILE_PATH, ROW_ID, SEQUENCE_NUMBER).map(_.toStructField))
 
   protected def createPaimonDataWriterFactory(info: PhysicalWriteInfo): 
DataWriterFactory = {
-    (_: Int, _: Long) =>
-      {
-        if (writeRowTracking) {
-          createPaimonMetadataAwareDataWriter()
-        } else {
-          PaimonV2DataWriter(
-            batchWriteBuilder,
-            writeSchema,
-            dataSchema,
-            coreOptions,
-            catalogContextForBlobDescriptor)
-        }
+    val uriReaderFactory = uriReaderFactoryForBlobDescriptor
+    (_: Int, _: Long) => {
+      if (writeRowTracking) {
+        createPaimonMetadataAwareDataWriter(uriReaderFactory)
+      } else {
+        PaimonV2DataWriter(
+          batchWriteBuilder,
+          writeSchema,
+          dataSchema,
+          coreOptions,
+          uriReaderFactory)
       }
+    }
   }
 
-  private def createPaimonMetadataAwareDataWriter(): PaimonV2DataWriter = {
+  private def createPaimonMetadataAwareDataWriter(
+      uriReaderFactory: UriReaderFactory): PaimonV2DataWriter = {
     new PaimonV2MetadataAwareDataWriter(
       batchWriteBuilder,
       writeSchema,
@@ -104,7 +106,7 @@ abstract class PaimonBatchWriteBase(
       dataSchema,
       rtMetadataSchema,
       coreOptions,
-      catalogContextForBlobDescriptor,
+      uriReaderFactory,
       rtPaimonWriteType)
   }
 
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
index 02b9338b4b..a954bd64b2 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDataWrite.scala
@@ -18,13 +18,13 @@
 
 package org.apache.paimon.spark.write
 
-import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.data.BinaryRow
 import org.apache.paimon.disk.IOManager
 import org.apache.paimon.spark.SparkUtils
 import org.apache.paimon.spark.util.SparkRowUtils
 import org.apache.paimon.table.sink._
 import org.apache.paimon.types.RowType
+import org.apache.paimon.utils.UriReaderFactory
 
 import org.apache.spark.sql.Row
 
@@ -37,7 +37,7 @@ case class PaimonDataWrite(
     writeRowTracking: Boolean = false,
     fullCompactionDeltaCommits: Option[Int],
     batchId: Option[Long],
-    catalogContext: CatalogContext,
+    uriReaderFactory: UriReaderFactory,
     postponePartitionBucketComputer: Option[BinaryRow => Integer])
   extends abstractInnerTableDataWrite[Row]
   with InnerTableV1DataWrite {
@@ -57,7 +57,7 @@ case class PaimonDataWrite(
   }
 
   private val toPaimonRow = {
-    SparkRowUtils.toPaimonRow(writeType, rowKindColIdx, catalogContext)
+    SparkRowUtils.toPaimonRow(writeType, rowKindColIdx, uriReaderFactory)
   }
 
   def write(row: Row): Unit = {
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriteBase.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriteBase.scala
index 3f1c6001fa..108649a294 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriteBase.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriteBase.scala
@@ -74,7 +74,7 @@ abstract class PaimonDeltaWriteBase(
       batchWriteBuilder,
       rowSchema,
       coreOptions,
-      catalogContextForBlobDescriptor,
+      uriReaderFactoryForBlobDescriptor,
       rowIdSchema.fieldIndex(FILE_PATH_COLUMN),
       rowIdSchema.fieldIndex(ROW_INDEX_COLUMN)
     )
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriter.scala
index 72bd73b846..7dacb329a6 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonDeltaWriter.scala
@@ -19,10 +19,10 @@
 package org.apache.paimon.spark.write
 
 import org.apache.paimon.CoreOptions
-import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.deletionvectors.{Bitmap64DeletionVector, 
BitmapDeletionVector, DeletionVector}
 import org.apache.paimon.spark.{PaimonDeletedRecordsTaskMetric, 
PaimonInsertedRecordsTaskMetric, PaimonUpdatedRecordsTaskMetric}
 import org.apache.paimon.table.sink.{BatchWriteBuilder, CommitMessage}
+import org.apache.paimon.utils.UriReaderFactory
 
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.connector.metric.CustomTaskMetric
@@ -52,7 +52,7 @@ case class PaimonDeltaWriterFactory(
     writeBuilder: BatchWriteBuilder,
     rowSchema: StructType,
     coreOptions: CoreOptions,
-    catalogContext: CatalogContext,
+    uriReaderFactory: UriReaderFactory,
     filePathOrdinal: Int,
     rowIndexOrdinal: Int)
   extends DeltaWriterFactory {
@@ -62,7 +62,7 @@ case class PaimonDeltaWriterFactory(
       writeBuilder,
       rowSchema,
       coreOptions,
-      catalogContext,
+      uriReaderFactory,
       filePathOrdinal,
       rowIndexOrdinal)
 }
@@ -79,7 +79,7 @@ case class PaimonDeltaWriter(
     writeBuilder: BatchWriteBuilder,
     rowSchema: StructType,
     coreOptions: CoreOptions,
-    catalogContext: CatalogContext,
+    uriReaderFactory: UriReaderFactory,
     filePathOrdinal: Int,
     rowIndexOrdinal: Int)
   extends DeltaWriter[InternalRow] {
@@ -96,7 +96,7 @@ case class PaimonDeltaWriter(
   private def getOrCreateAppendWriter: PaimonV2DataWriter = {
     appendWriter.getOrElse {
       val writer =
-        PaimonV2DataWriter(writeBuilder, rowSchema, rowSchema, coreOptions, 
catalogContext)
+        PaimonV2DataWriter(writeBuilder, rowSchema, rowSchema, coreOptions, 
uriReaderFactory)
       appendWriter = Some(writer)
       writer
     }
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonV2DataWriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonV2DataWriter.scala
index eadd056cf2..05b6bfd307 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonV2DataWriter.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/PaimonV2DataWriter.scala
@@ -19,12 +19,11 @@
 package org.apache.paimon.spark.write
 
 import org.apache.paimon.CoreOptions
-import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.spark.{SparkInternalRowWrapper, SparkUtils}
 import org.apache.paimon.spark.metric.SparkMetricRegistry
 import org.apache.paimon.table.sink.{BatchWriteBuilder, CommitMessage, 
TableWriteImpl}
 import org.apache.paimon.types.RowType
-import org.apache.paimon.utils.IOUtils
+import org.apache.paimon.utils.{IOUtils, UriReaderFactory}
 
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.catalyst.expressions.JoinedRow
@@ -38,7 +37,7 @@ case class PaimonV2DataWriter(
     writeSchema: StructType,
     dataSchema: StructType,
     coreOptions: CoreOptions,
-    catalogContext: CatalogContext,
+    uriReaderFactory: UriReaderFactory,
     batchId: Option[Long] = None,
     paimonWriteType: Option[RowType] = None,
     metadataSchema: Option[StructType] = None,
@@ -79,7 +78,7 @@ case class PaimonV2DataWriter(
       schema: StructType): InternalRow => SparkInternalRowWrapper = {
     val numFields = writeSchema.fields.length
     val reusableWrapper =
-      new SparkInternalRowWrapper(writeSchema, numFields, schema, 
catalogContext)
+      SparkInternalRowWrapper.fromUriReaderFactory(writeSchema, numFields, 
schema, uriReaderFactory)
     record => reusableWrapper.replace(record)
   }
 
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/WriteHelper.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/WriteHelper.scala
index c039bf80ed..9de665c12b 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/WriteHelper.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/write/WriteHelper.scala
@@ -20,13 +20,12 @@ package org.apache.paimon.spark.write
 
 import org.apache.paimon.CoreOptions
 import org.apache.paimon.CoreOptions.TagCreationMode
-import org.apache.paimon.catalog.CatalogContext
 import org.apache.paimon.partition.actions.PartitionMarkDoneAction
 import org.apache.paimon.spark.catalyst.Compatibility
-import org.apache.paimon.table.FileStoreTable
+import org.apache.paimon.table.{BlobDescriptorReaderFactory, FileStoreTable}
 import org.apache.paimon.table.sink.CommitMessage
 import org.apache.paimon.tag.TagBatchCreation
-import org.apache.paimon.utils.{BlobDescriptorUtils, 
InternalRowPartitionComputer, PartitionPathUtils, PartitionStatisticsReporter, 
TypeUtils}
+import org.apache.paimon.utils.{InternalRowPartitionComputer, 
PartitionPathUtils, PartitionStatisticsReporter, TypeUtils, UriReaderFactory}
 
 import org.apache.spark.internal.Logging
 import org.apache.spark.sql.PaimonSparkSession
@@ -42,10 +41,8 @@ trait WriteHelper extends Logging {
 
   lazy val coreOptions: CoreOptions = table.coreOptions()
 
-  lazy val catalogContextForBlobDescriptor: CatalogContext =
-    BlobDescriptorUtils.getCatalogContext(
-      table.catalogEnvironment().catalogContext(),
-      coreOptions.toConfiguration)
+  lazy val uriReaderFactoryForBlobDescriptor: UriReaderFactory =
+    BlobDescriptorReaderFactory.create(table)
 
   // Spark support v2 write driver metrics since 4.0, see 
https://github.com/apache/spark/pull/48573
   // To ensure compatibility with 3.x, manually post driver metrics here 
instead of using Spark's API.
diff --git a/paimon-spark/paimon-spark-ut/pom.xml 
b/paimon-spark/paimon-spark-ut/pom.xml
index ce8fae9b2f..41fa5a86f5 100644
--- a/paimon-spark/paimon-spark-ut/pom.xml
+++ b/paimon-spark/paimon-spark-ut/pom.xml
@@ -95,6 +95,13 @@ under the License.
             <version>${project.version}</version>
             <scope>test</scope>
         </dependency>
+
+        <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <version>${mockito.version}</version>
+            <scope>test</scope>
+        </dependency>
     </dependencies>
 
     <build>
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
 
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
index 1117ad58c7..c7ecc35b43 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
+++ 
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
@@ -19,21 +19,29 @@
 package org.apache.paimon.spark;
 
 import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.BlobDescriptor;
 import org.apache.paimon.data.Decimal;
 import org.apache.paimon.data.GenericArray;
 import org.apache.paimon.data.GenericMap;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.fs.local.LocalFileIO;
 import org.apache.paimon.spark.data.SparkInternalRow;
 import org.apache.paimon.utils.DateTimeUtils;
+import org.apache.paimon.utils.UriReaderFactory;
 
 import org.apache.spark.sql.catalyst.CatalystTypeConverters;
+import org.apache.spark.sql.catalyst.expressions.GenericInternalRow;
 import org.apache.spark.sql.catalyst.util.CharVarcharUtils;
+import org.apache.spark.sql.types.DataTypes;
+import org.apache.spark.sql.types.StructType;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
 
 import java.math.BigDecimal;
 import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
 import java.time.LocalDate;
 import java.time.LocalDateTime;
 import java.util.AbstractMap;
@@ -52,6 +60,8 @@ import static org.assertj.core.api.Assertions.assertThat;
 /** Test for {@link SparkInternalRow}. */
 public class SparkInternalRowTest {
 
+    @TempDir java.nio.file.Path tempPath;
+
     @Test
     public void test() {
         TimeZone tz = TimeZone.getDefault();
@@ -129,6 +139,26 @@ public class SparkInternalRowTest {
         TimeZone.setDefault(tz);
     }
 
+    @Test
+    public void testReadBlobWithProvidedUriReaderFactory() throws Exception {
+        byte[] bytes = new byte[] {1, 2, 3};
+        java.nio.file.Path blobFile = tempPath.resolve("blob");
+        Files.write(blobFile, bytes);
+        byte[] descriptor =
+                new BlobDescriptor(blobFile.toUri().toString(), 0, 
bytes.length).serialize();
+        StructType schema = new StructType().add("blob", DataTypes.BinaryType);
+
+        SparkInternalRowWrapper wrapper =
+                SparkInternalRowWrapper.fromUriReaderFactory(
+                                schema,
+                                1,
+                                schema,
+                                
UriReaderFactory.fromFileIO(LocalFileIO.create()))
+                        .replace(new GenericInternalRow(new Object[] 
{descriptor}));
+
+        assertThat(wrapper.getBlob(0).toData()).isEqualTo(bytes);
+    }
+
     private String sparkRowToString(org.apache.spark.sql.Row row) {
         return JavaConverters.seqAsJavaList(row.toSeq()).stream()
                 .map(
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/write/PaimonBatchWriteTest.java
 
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/write/PaimonBatchWriteTest.java
new file mode 100644
index 0000000000..6db8176d63
--- /dev/null
+++ 
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/write/PaimonBatchWriteTest.java
@@ -0,0 +1,82 @@
+/*
+ * 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.paimon.spark.write;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.CatalogLoader;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.table.CatalogEnvironment;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+
+import org.apache.spark.sql.types.StructType;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import scala.Option;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests for {@link PaimonBatchWrite}. */
+public class PaimonBatchWriteTest {
+
+    @Test
+    public void testResolveBlobDescriptorSourceBeforeCreatingTaskWriters() 
throws Exception {
+        Identifier sourceIdentifier = Identifier.fromString("db.source");
+        Table sourceTable = mock(Table.class);
+        when(sourceTable.fileIO()).thenReturn(LocalFileIO.create());
+
+        Catalog catalog = mock(Catalog.class);
+        when(catalog.getTable(sourceIdentifier)).thenReturn(sourceTable);
+        CatalogLoader catalogLoader = mock(CatalogLoader.class);
+        when(catalogLoader.load()).thenReturn(catalog);
+        CatalogEnvironment catalogEnvironment = mock(CatalogEnvironment.class);
+        when(catalogEnvironment.catalogLoader()).thenReturn(catalogLoader);
+
+        FileStoreTable targetTable = mock(FileStoreTable.class);
+        when(targetTable.catalogEnvironment()).thenReturn(catalogEnvironment);
+        when(targetTable.coreOptions())
+                .thenReturn(
+                        CoreOptions.fromMap(
+                                Collections.singletonMap(
+                                        "blob-descriptor.source-table", 
"db.source")));
+        
when(targetTable.newBatchWriteBuilder()).thenReturn(mock(BatchWriteBuilder.class));
+
+        StructType schema = new StructType();
+        PaimonBatchWrite batchWrite =
+                new PaimonBatchWrite(
+                        targetTable,
+                        schema,
+                        schema,
+                        Option.empty(),
+                        Option.empty(),
+                        Option.empty());
+
+        batchWrite.createBatchWriterFactory(null);
+
+        verify(catalogLoader).load();
+        verify(catalog).getTable(sourceIdentifier);
+    }
+}
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala
index feb8acfbfb..eaa122b964 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/BlobTestBase.scala
@@ -243,6 +243,36 @@ class BlobTestBase extends PaimonSparkTestBase {
     }
   }
 
+  test("Blob: materialize descriptor with source table FileIO") {
+    withTable("blob_source", "blob_target") {
+      sql(
+        "CREATE TABLE blob_source (id INT, picture BINARY) TBLPROPERTIES (" +
+          "'row-tracking.enabled'='true', " +
+          "'data-evolution.enabled'='true', " +
+          "'blob-field'='picture')")
+      sql("INSERT INTO blob_source VALUES (1, X'48656C6C6F'), (2, X'5945')")
+      sql(
+        "ALTER TABLE blob_source SET TBLPROPERTIES (" +
+          "'blob-as-descriptor'='true')")
+
+      sql(
+        "CREATE TABLE blob_target (id INT, picture BINARY) TBLPROPERTIES (" +
+          "'row-tracking.enabled'='true', " +
+          "'data-evolution.enabled'='true', " +
+          "'blob-field'='picture', " +
+          s"'blob-descriptor.source-table'='$dbName0.blob_source')")
+      sql("INSERT INTO blob_target SELECT * FROM blob_source")
+
+      checkAnswer(
+        sql("SELECT id, picture FROM blob_target ORDER BY id"),
+        Seq(
+          Row(1, Array[Byte](72, 101, 108, 108, 111)),
+          Row(2, Array[Byte](89, 69))
+        )
+      )
+    }
+  }
+
   test("Blob: test write blob descriptor with partition") {
     withTable("t") {
       val blobData = new Array[Byte](1024 * 1024)
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
index 1bc3cc1888..72ca496a12 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/PostponeBucketTableTest.scala
@@ -18,9 +18,13 @@
 
 package org.apache.paimon.spark.sql
 
+import org.apache.paimon.catalog.{Catalog, CatalogLoader, DelegateCatalog, 
Identifier}
 import org.apache.paimon.fs.Path
 import org.apache.paimon.spark.PaimonSparkTestBase
+import org.apache.paimon.spark.procedure.SparkPostponeCompactProcedure
+import org.apache.paimon.table.{CatalogEnvironment, FileStoreTableFactory}
 
+import org.apache.spark.TaskContext
 import org.apache.spark.sql.Row
 
 class PostponeBucketTableTest extends PaimonSparkTestBase {
@@ -230,6 +234,46 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
     }
   }
 
+  test("Postpone compaction resolves blob descriptor source outside Spark 
tasks") {
+    withTable("source", "t") {
+      sql("CREATE TABLE source (k INT, v STRING) TBLPROPERTIES ('primary-key' 
= 'k')")
+      sql("""
+            |CREATE TABLE t (
+            |  k INT,
+            |  v STRING
+            |) TBLPROPERTIES (
+            |  'primary-key' = 'k',
+            |  'bucket' = '-2',
+            |  'postpone.batch-write-fixed-bucket' = 'false',
+            |  'blob-descriptor.source-table' = 'test.source'
+            |)
+            |""".stripMargin)
+      sql("INSERT INTO t SELECT id, CAST(id AS STRING) FROM range(0, 20)")
+
+      val table = loadTable("t")
+      val environment = table.catalogEnvironment
+      val driverOnlyEnvironment = new CatalogEnvironment(
+        environment.identifier,
+        environment.uuid,
+        new 
PostponeBucketTableTest.SourceTableDriverOnlyCatalogLoader(environment.catalogLoader),
+        environment.lockFactory,
+        environment.lockContext,
+        environment.catalogContext,
+        environment.supportsVersionManagement,
+        false
+      )
+      val driverOnlyTable = FileStoreTableFactory.create(
+        table.fileIO,
+        table.location,
+        table.schema,
+        driverOnlyEnvironment)
+
+      SparkPostponeCompactProcedure(driverOnlyTable, spark, null, 
createRelationV2("t")).execute()
+
+      checkAnswer(sql("SELECT count(*) FROM t"), Seq(Row(20)))
+    }
+  }
+
   test("Postpone partition bucket table: write postpone bucket then compact") {
     withTable("t") {
       sql("""
@@ -375,3 +419,22 @@ class PostponeBucketTableTest extends PaimonSparkTestBase {
     }
   }
 }
+
+object PostponeBucketTableTest {
+
+  private class SourceTableDriverOnlyCatalogLoader(delegate: CatalogLoader) 
extends CatalogLoader {
+
+    override def load(): Catalog = {
+      new DelegateCatalog(delegate.load()) {
+        override def catalogLoader(): CatalogLoader = delegate
+
+        override def getTable(identifier: Identifier) = {
+          if (TaskContext.get() != null && identifier.getTableName == 
"source") {
+            throw new IllegalStateException("Source table must not be loaded 
in a Spark task")
+          }
+          super.getTable(identifier)
+        }
+      }
+    }
+  }
+}

Reply via email to