[bugfix] forward shuffle config to odps and kafka readers - #683
Conversation
BaseReader implements shuffle as a batch buffer driven by its shuffle and shuffle_buffer_size arguments, but OdpsDataset never passed the two data_config fields when constructing OdpsReader, and KafkaReader neither accepted them nor passed them to BaseReader. Both readers therefore kept the shuffle=False default and data_config.shuffle was silently ignored. Both datasets now forward the fields the way the csv and parquet datasets do, enabling shuffle only in train mode, and KafkaReader relays them to BaseReader. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RwaumtMBn6cQvR51aJjAPU
068f9d0 to
e15bf77
Compare
| self._batch_size, | ||
| list(self._selected_input_names) if self._selected_input_names else None, | ||
| self._data_config.drop_remainder, | ||
| shuffle=self._data_config.shuffle and self._mode == Mode.TRAIN, |
There was a problem hiding this comment.
The forwarding mechanics here are correct (verified against BaseReader.__init__ and the csv/parquet/odps_v1 call sites). The one thing worth an explicit signal to users is that the PR body's Parquet-parity framing understates the Kafka case, in two Kafka-specific ways:
-
Resume offset loss is permanent, and resume is the normal operating mode.
update_dataloder_statemergesCKPT_ROW_IDX(= message offset,kafka_dataset.py:586) by max over yielded batches, andon_assignseeks tocheckpoint_state[key] + 1(kafka_dataset.py:417-419). Without shuffle, offsets are monotonic so the max equals the true consumed prefix and resume is lossless. With shuffle on, up toshuffle_buffer_size - 1already-consumed-but-buffered batches below the recorded max are skipped forever at each failover. Parquet/ODPS are finite passes where a later epoch re-covers the skipped window; the Kafka consume loop never ends, so nothing ever re-covers it — and this PR is the first time shuffle can be on for that path. -
The event-time watermark loses its lower-bound property.
Batch.data_timestampis the max event-time of the yielded batch (dataset.py:370-375), stamped asDATA_TS_WATERMARKand used to drive event-time checkpoint triggers. Out-of-order pops mean a checkpoint can record watermarkTwhile messages withts < Tare still in the buffer (and then skipped per point 1), and a late-popped old batch can regress the observed timestamp, skewingshould_save_on_timestamptriggers.
Neither needs to block this fix, but shuffle=true on a Kafka source (especially ODL with event-time checkpointing) now carries a non-obvious correctness tradeoff with no user-visible signal. Suggestion: a logger.warning when KafkaReader is constructed with shuffle enabled, and/or a caveat in the Kafka docs.
Review summaryVerdict: the fix is correct and well-tested. Independently verified:
One noteworthy concern, posted inline on Minor notes:
All five review areas (code quality, performance, test coverage, documentation accuracy, security) completed. 🤖 Generated with Claude Code |
BaseReaderalready implementsdata_config.shuffleas a buffer ofshuffle_buffer_sizebatches that is drained in random order, andOdpsReaderhas accepted theshuffle/shuffle_buffer_sizearguments since the feature landed. ButOdpsDatasetnever passed them when constructing its reader, andKafkaReaderneither accepted them nor relayed them toBaseReader. Both readers therefore kept theshuffle=Falsedefault, so settingdata_config.shufflehad no effect onOdpsDatasetorKafkaDatasetwhile it worked onCsvDataset/ParquetDataset.Both datasets now forward the two fields exactly like the csv and parquet datasets, gated on train mode, and
KafkaReaderrelays them toBaseReader.Note that the shuffle buffer reorders batches, so with shuffle on, the per-source max row index stored in a dataloader checkpoint can run ahead of batches still sitting in the buffer, and a resume skips them. That tradeoff already exists for
ParquetDatasetand is unchanged here.Test Plan
Added a broker/ODPS-free regression test to each dataset's test module asserting that the reader receives
shuffle=Trueandshuffle_buffer_size=64inMode.TRAINandshuffle=FalseinMode.EVAL. Both fail on master (KeyError: 'shuffle'for odps,False != Truefor kafka) and pass with the fix.The assertion is made on the reader construction rather than on observed batch order because a real reordering check needs a live ODPS table or kafka broker, which the existing integration tests in these modules already gate behind CI credentials, and those tests run in
Mode.EVALwhere shuffle stays off by design.🤖 Generated with Claude Code
https://claude.ai/code/session_01RwaumtMBn6cQvR51aJjAPU