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 dbb5423628 [Bug] Fix concurrent LocalTableQuery lookups returning null
across files (#9500)
dbb5423628 is described below
commit dbb5423628935752d0533e50361eb2aea9b9d1ed
Author: Arvin <[email protected]>
AuthorDate: Wed Sep 2 14:43:28 2026 +0800
[Bug] Fix concurrent LocalTableQuery lookups returning null across files
(#9500)
---
.../data/serializer/RowCompactedSerializer.java | 25 ++++++++++++----------
1 file changed, 14 insertions(+), 11 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/serializer/RowCompactedSerializer.java
b/paimon-common/src/main/java/org/apache/paimon/data/serializer/RowCompactedSerializer.java
index baa4d04325..c2a501b7e3 100644
---
a/paimon-common/src/main/java/org/apache/paimon/data/serializer/RowCompactedSerializer.java
+++
b/paimon-common/src/main/java/org/apache/paimon/data/serializer/RowCompactedSerializer.java
@@ -735,14 +735,15 @@ public class RowCompactedSerializer implements
Serializer<InternalRow> {
private static class SliceComparator implements Comparator<MemorySlice> {
- private final RowReader reader1;
- private final RowReader reader2;
+ private final int headerSizeInBytes;
+ private final ThreadLocal<RowReader> reader1;
+ private final ThreadLocal<RowReader> reader2;
private final FieldReader[] fieldReaders;
public SliceComparator(RowType rowType) {
- int bitSetInBytes =
calculateBitSetInBytes(rowType.getFieldCount());
- this.reader1 = new RowReader(bitSetInBytes);
- this.reader2 = new RowReader(bitSetInBytes);
+ this.headerSizeInBytes =
calculateBitSetInBytes(rowType.getFieldCount());
+ this.reader1 = ThreadLocal.withInitial(() -> new
RowReader(headerSizeInBytes));
+ this.reader2 = ThreadLocal.withInitial(() -> new
RowReader(headerSizeInBytes));
this.fieldReaders = new FieldReader[rowType.getFieldCount()];
for (int i = 0; i < rowType.getFieldCount(); i++) {
fieldReaders[i] = createFieldReader(rowType.getTypeAt(i));
@@ -751,11 +752,13 @@ public class RowCompactedSerializer implements
Serializer<InternalRow> {
@Override
public int compare(MemorySlice slice1, MemorySlice slice2) {
- reader1.pointTo(slice1.segment(), slice1.offset());
- reader2.pointTo(slice2.segment(), slice2.offset());
+ RowReader r1 = reader1.get();
+ RowReader r2 = reader2.get();
+ r1.pointTo(slice1.segment(), slice1.offset());
+ r2.pointTo(slice2.segment(), slice2.offset());
for (int i = 0; i < fieldReaders.length; i++) {
- boolean isNull1 = reader1.isNullAt(i);
- boolean isNull2 = reader2.isNullAt(i);
+ boolean isNull1 = r1.isNullAt(i);
+ boolean isNull2 = r2.isNullAt(i);
if (!isNull1 || !isNull2) {
if (isNull1) {
return -1;
@@ -763,8 +766,8 @@ public class RowCompactedSerializer implements
Serializer<InternalRow> {
return 1;
} else {
FieldReader fieldReader = fieldReaders[i];
- Object o1 = fieldReader.readField(reader1, i);
- Object o2 = fieldReader.readField(reader2, i);
+ Object o1 = fieldReader.readField(r1, i);
+ Object o2 = fieldReader.readField(r2, i);
int comp;
if (o1 instanceof byte[]) {
// BINARY / VARBINARY fields read back as byte[],
which does not