Skip to content

Commit d373e0d

Browse files
committed
fix wrapper
1 parent e83a839 commit d373e0d

1 file changed

Lines changed: 5 additions & 8 deletions

File tree

fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumerator.java

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,6 @@
5050
import org.apache.fluss.metadata.TableInfo;
5151
import org.apache.fluss.metadata.TablePath;
5252
import org.apache.fluss.predicate.Predicate;
53-
import org.apache.fluss.row.InternalRow;
5453
import org.apache.fluss.shaded.guava32.com.google.common.collect.Lists;
5554
import org.apache.fluss.types.RowType;
5655
import org.apache.fluss.utils.ExceptionUtils;
@@ -564,7 +563,11 @@ private List<PartitionInfo> applyPartitionFilter(List<PartitionInfo> partitionIn
564563
.filter(
565564
partition ->
566565
partitionFilters.test(
567-
toInternalRow(partition, partitionRowType)))
566+
PartitionUtils.toPartitionRow(
567+
partition
568+
.getResolvedPartitionSpec()
569+
.getPartitionValues(),
570+
partitionRowType)))
568571
.collect(Collectors.toList());
569572

570573
int filteredSize = filteredPartitionInfos.size();
@@ -587,12 +590,6 @@ private List<PartitionInfo> applyPartitionFilter(List<PartitionInfo> partitionIn
587590
}
588591
}
589592

590-
private static InternalRow toInternalRow(
591-
PartitionInfo partitionInfo, RowType partitionRowType) {
592-
return PartitionUtils.toPartitionRow(
593-
partitionInfo.getResolvedPartitionSpec().getPartitionValues(), partitionRowType);
594-
}
595-
596593
/** Init the splits for Fluss. */
597594
private void checkPartitionChanges(Set<PartitionInfo> partitionInfos, Throwable t) {
598595
if (closed) {

0 commit comments

Comments
 (0)