I am facing issue in using StreamingQueryListener to identify number of input row, I am using
queryProgress.progress().numInputRows()
I get proper count when there is no other action apart from write, but the moment I add certain actions such as df.count or df.isEmpty() my number of input rows count get disrupted.
Any help is highly appreciated
EDIT
Below code works
df.writeStream().outputMode("append").foreachBatch(new VoidFunction2<Dataset<Row>,Long>(){
@Override
public void call(Dataset<Row> streamDataset, Long batchId) throws Exception {
streamDataset.write().mode(SaveMode.Append).save("namesAndFavColors.parquet");
}
}).start();
This gives wrong count
df.writeStream().outputMode("append").foreachBatch(new VoidFunction2<Dataset<Row>,Long>(){
@Override
public void call(Dataset<Row> streamDataset, Long batchId) throws Exception {
streamDataset.count();
streamDataset.write().mode(SaveMode.Append).save("namesAndFavColors.parquet");
}
}).start();
Note
Please ignore write() code, in real scenario data is being written to mysql