codope commented on code in PR #9943:
URL: https://github.com/apache/hudi/pull/9943#discussion_r1375791876
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/helpers/QueryRunner.java:
##########
@@ -80,22 +81,30 @@ public static Dataset<Row> applyOrdering(Dataset<Row>
dataset, List<String> orde
return dataset;
}
- public Dataset<Row> runIncrementalQuery(QueryInfo queryInfo) {
+ public Pair<QueryInfo, Dataset<Row>> runIncrementalQuery(QueryInfo
queryInfo) {
LOG.info("Running incremental query");
- return sparkSession.read().format("org.apache.hudi")
+ return Pair.of(queryInfo, sparkSession.read().format("org.apache.hudi")
.option(DataSourceReadOptions.QUERY_TYPE().key(),
queryInfo.getQueryType())
.option(DataSourceReadOptions.BEGIN_INSTANTTIME().key(),
queryInfo.getPreviousInstant())
- .option(DataSourceReadOptions.END_INSTANTTIME().key(),
queryInfo.getEndInstant()).load(sourcePath);
+ .option(DataSourceReadOptions.END_INSTANTTIME().key(),
queryInfo.getEndInstant()).load(sourcePath));
}
- public Dataset<Row> runSnapshotQuery(QueryInfo queryInfo) {
+ public Pair<QueryInfo, Dataset<Row>> runSnapshotQuery(QueryInfo queryInfo,
Option<SnapshotLoadQuerySplitter> snapshotLoadQuerySplitterOption) {
LOG.info("Running snapshot query");
- return sparkSession.read().format("org.apache.hudi")
- .option(DataSourceReadOptions.QUERY_TYPE().key(),
queryInfo.getQueryType()).load(sourcePath)
+ Dataset<Row> snapshot = sparkSession.read().format("org.apache.hudi")
+ .option(DataSourceReadOptions.QUERY_TYPE().key(),
queryInfo.getQueryType()).load(sourcePath);
+ QueryInfo snapshotQueryInfo = snapshotLoadQuerySplitterOption
+ .map(snapshotLoadQuerySplitter ->
snapshotLoadQuerySplitter.getNextCheckpoint(snapshot, queryInfo))
Review Comment:
If you do `Option.ofNullable` and the snapshot query load splitter class was
not set then this could lead to NPE right? Please check my comment above.
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/GcsEventsHoodieIncrSource.java:
##########
@@ -145,6 +148,9 @@ public GcsEventsHoodieIncrSource(TypedProperties props,
JavaSparkContext jsc, Sp
this.gcsObjectDataFetcher = gcsObjectDataFetcher;
this.queryRunner = queryRunner;
this.schemaProvider = Option.ofNullable(schemaProvider);
+ this.snapshotLoadQuerySplitter =
Option.ofNullable(props.getString(SNAPSHOT_LOAD_QUERY_SPLITTER_CLASS_NAME,
null))
Review Comment:
Not setting `SNAPSHOT_LOAD_QUERY_SPLITTER_CLASS_NAME` could lead to NPE
somewhere. Maybe just validate that the config is set and if not then
Option.empty instead of Option.ofNullable. Then, wherever it is used handle
empty if necessary.
--
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]