I am writing a custom structured streaming source and can't figure out how batch sizes can be limited. The new MicroBatchStream interface provides the planInputPartitions method that gets called with the return value of latestOffset as end Offset. It then returns a partitioning of the data up to the provided latest offset to be processed in a single batch.
When I start a new streaming query this results in an enormously large first batch as all historic data gets crammed into a single batch.
I already tried manually limiting the batch size by gradually increasing the latestOffset depending on what has already been committed. When restarting a query from a checkpoint, however, this fails as nothing has been committed yet.
Is there a (obvious) way of limiting streaming batch sizes?