Skip to content

[bugfix] forward shuffle config to odps and kafka readers - #683

Merged
tiankongdeguiji merged 1 commit into
masterfrom
bugfix/shuffle-odps-kafka-dataset
Sep 24, 2026
Merged

tiankongdeguiji merged 1 commit into
masterfrom
bugfix/shuffle-odps-kafka-dataset

Conversation

@tiankongdeguiji

Copy link
Copy Markdown
Collaborator

BaseReader already implements data_config.shuffle as a buffer of shuffle_buffer_size batches that is drained in random order, and OdpsReader has accepted the shuffle / shuffle_buffer_size arguments since the feature landed. But OdpsDataset never passed them when constructing its reader, and KafkaReader neither accepted them nor relayed them to BaseReader. Both readers therefore kept the shuffle=False default, so setting data_config.shuffle had no effect on OdpsDataset or KafkaDataset while it worked on CsvDataset / ParquetDataset.

Both datasets now forward the two fields exactly like the csv and parquet datasets, gated on train mode, and KafkaReader relays them to BaseReader.

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 ParquetDataset and is unchanged here.

Test Plan

Added a broker/ODPS-free regression test to each dataset's test module asserting that the reader receives shuffle=True and shuffle_buffer_size=64 in Mode.TRAIN and shuffle=False in Mode.EVAL. Both fail on master (KeyError: 'shuffle' for odps, False != True for 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.EVAL where shuffle stays off by design.

python -m unittest tzrec.datasets.odps_dataset_test tzrec.datasets.kafka_dataset_test tzrec.datasets.dataset_test
# Ran 78 tests, OK (skipped=22, the credentialed odps/kafka integration tests)
python -m unittest tzrec.datasets.csv_dataset_test tzrec.datasets.parquet_dataset_test
# Ran 30 tests, OK
pre-commit run --files <changed files>   # all passed
pyrefly check                            # 0 errors

🤖 Generated with Claude Code

https://claude.ai/code/session_01RwaumtMBn6cQvR51aJjAPU

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
@tiankongdeguiji
tiankongdeguiji force-pushed the bugfix/shuffle-odps-kafka-dataset branch from 068f9d0 to e15bf77 Compare September 21, 2026 07:25
@tiankongdeguiji tiankongdeguiji added the claude-review Let Claude Review label Sep 21, 2026
@github-actions github-actions Bot removed the claude-review Let Claude Review label Sep 21, 2026
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,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  1. Resume offset loss is permanent, and resume is the normal operating mode. update_dataloder_state merges CKPT_ROW_IDX (= message offset, kafka_dataset.py:586) by max over yielded batches, and on_assign seeks to checkpoint_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 to shuffle_buffer_size - 1 already-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.

  2. The event-time watermark loses its lower-bound property. Batch.data_timestamp is the max event-time of the yielded batch (dataset.py:370-375), stamped as DATA_TS_WATERMARK and used to drive event-time checkpoint triggers. Out-of-order pops mean a checkpoint can record watermark T while messages with ts < T are still in the buffer (and then skipped per point 1), and a late-popped old batch can regress the observed timestamp, skewing should_save_on_timestamp triggers.

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.

@github-actions

Copy link
Copy Markdown
Contributor

Review summary

Verdict: the fix is correct and well-tested. Independently verified:

  • The shuffle=data_config.shuffle and mode == Mode.TRAIN gate and kwarg forwarding match the CsvDataset / ParquetDataset / odps_dataset_v1 call sites exactly; KafkaReader's positional relay matches BaseReader.__init__'s parameter order (no argument swap); both readers route to_batches through BaseReader._arrow_reader_iter, so shuffle genuinely takes effect.
  • The new tests pin the bug (both fail on master as described), are broker/ODPS-free — the Kafka one avoids _peek_schema_from_kafka because input_fields is populated, the ODPS one patches the module-level OdpsReader — and need no CI-scope marker.
  • Docs review found nothing stale: docs/source/feature/data.md describes shuffle generically and becomes true for odps/kafka with this fix. Security review found no rank-disjointness or seeding issues (shuffle runs strictly after per-worker slicing), and shuffle_buffer_size is proto2 uint32 with a default matching BaseReader's.

One noteworthy concern, posted inline on kafka_dataset.py:250: the Parquet-parity framing in the PR body understates the Kafka case — on an endless stream, max-offset checkpointing plus a shuffle buffer turns failover resume into silent permanent message skips, and the Kafka-only event-time watermark (DATA_TS_WATERMARK) loses its lower-bound property. Suggested a warning log and/or docs caveat there; not necessarily blocking.

Minor notes:

  • After merge, existing odps/kafka configs that already set shuffle=true will start actually shuffling and paying the per-worker shuffle-buffer memory — worth a release-note line.
  • Optional follow-up (not for this bugfix): the TRAIN gate is now duplicated in four dataset constructors; centralizing it in BaseDataset.__init__ would prevent the same omission in the next dataset.

All five review areas (code quality, performance, test coverage, documentation accuracy, security) completed.

🤖 Generated with Claude Code

@tiankongdeguiji
tiankongdeguiji merged commit d0e3ce5 into master Sep 24, 2026
12 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants