ruanhang1993 commented on code in PR #21589:
URL: https://github.com/apache/flink/pull/21589#discussion_r1080938338
##########
flink-connectors/flink-connector-base/src/main/java/org/apache/flink/connector/base/source/reader/SourceReaderBase.java:
##########
@@ -328,14 +348,82 @@ private SplitContext(String splitId, SplitStateT state) {
this.splitId = splitId;
}
- SourceOutput<T> getOrCreateSplitOutput(ReaderOutput<T> mainOutput) {
+ SourceOutput<T> getOrCreateSplitOutput(
+ ReaderOutput<T> mainOutput,
+ @Nullable RecordEvaluator<T> recordEvaluator,
+ SplitT split,
+ SplitFetcherManager<E, SplitT> splitFetcherManager) {
if (sourceOutput == null) {
// The split output should have been created when
AddSplitsEvent was processed in
// SourceOperator. Here we just use this method to get the
previously created
// output.
sourceOutput = mainOutput.createOutputForSplit(splitId);
+ if (recordEvaluator != null) {
+ sourceOutput =
+ new SourceOutputWrapper<>(
+ split, recordEvaluator, sourceOutput,
splitFetcherManager);
+ }
}
return sourceOutput;
}
}
+
+ private static final class SourceOutputWrapper<E, T, SplitT extends
SourceSplit>
+ implements SourceOutput<T> {
+ final SplitT split;
+ final RecordEvaluator<T> recordEvaluator;
+ final SourceOutput<T> sourceOutput;
+ final SplitFetcherManager<E, SplitT> splitFetcherManager;
+
+ private boolean isStreamEnd = false;
+
+ public SourceOutputWrapper(
+ SplitT split,
+ RecordEvaluator<T> recordEvaluator,
+ SourceOutput<T> sourceOutput,
+ SplitFetcherManager<E, SplitT> splitFetcherManager) {
+ this.split = split;
+ this.recordEvaluator = recordEvaluator;
+ this.sourceOutput = sourceOutput;
+ this.splitFetcherManager = splitFetcherManager;
+ }
+
+ @Override
+ public void emitWatermark(Watermark watermark) {
+ sourceOutput.emitWatermark(watermark);
+ }
+
+ @Override
+ public void markIdle() {
+ sourceOutput.markIdle();
+ }
+
+ @Override
+ public void markActive() {
+ sourceOutput.markActive();
+ }
+
+ @Override
+ public void collect(T record) {
+ if (!isStreamEnd) {
+ handleEndOfStreamRecord(record);
Review Comment:
I forget that the end record will not be sent. I will fix it.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]