Repository navigation
Nemotron OCR SDG Pipeline - #1899
Conversation
- Updated the `transformers` dependency from `<=4.55.2` to `==4.57.0` in `pyproject.toml` and `uv.lock` to ensure compatibility with the Cosmos Embed imports. - Added a new `Gemini3Pro` model class in `gemini.py` utilizing the NVIDIA Inference API. - Introduced `DescriptionOutputStage` and `DescriptionValidatorStage` for processing and validating image descriptions, respectively. - Enhanced `VLMProcessingStage` to improve GPU resource handling and added a `num_workers` parameter to `DescriptionStage` for better scalability. This commit enhances the model's capabilities and ensures that dependencies are up-to-date for optimal performance. Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Removes description-specific stages (description*.py, description pipeline tutorials) that belong on aot/omni_description. Adds OCR result inspection/review scripts and shared design docs. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
… pipeline tutorial Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
…Data and OCRDenseWord classes Signed-off-by: Ao Tang <aot@nvidia.com>
--metrics-dir wires Ray metrics into the running Prometheus/Grafana instance. --run-name sets SLURM_JOB_NAME so Xenna labels the run on the ray_pipeline_input_tasks metric for human-readable identification in Grafana. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
… timing - OCRConversationData.to_dict(): call conversation.to_dict() explicitly instead of relying on dataclasses.asdict(), which bypasses the custom ConversationSample serialization and drops the "t" media-type field from image fragments. - RayClient: move Prometheus service-discovery registration to after Ray is started and responsive; add _wait_for_ray_service_discovery_file() so the SD file exists before Prometheus is told to watch it. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
Remove internal-only I/O stages from io.py (InputFormat, ImageReaderStage, ImageFolderReaderStage, TarImageReaderStage, JsonlTarImageReaderStage, OcrJsonlReaderStage, JsonlPipelineOutputReaderStage, TarImageReader, ParquetImageReaderStage, ParquetImageReader). These classes are only used by the internal ocr_pipeline.py and will live on aot/omni_sdg_internal. Public io.py now exports: HFDatasetImageReaderStage, SkipProcessedStage, ResultWriterStage, merge_output_shards, ImageWriterStage, and the FileReader helpers (load_image_from_task, TarFileReader, etc.). Add tests/stages/test_hf_dataset_image_reader.py covering HFDatasetImageReaderStage. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
ImageWriterStage is not used by hf_ocr_pipeline.py and only makes sense with the internal JSONL+tar pipeline; moved to aot/omni_sdg_internal. TarFileReader, ParquetFileReader, _file_readers dispatcher, and the deprecated _parse_tar_slice_path wrapper are removed. In the public HF pipeline images are always regular JPEG files on disk, so load_image_from_task is simplified to a single RegularFileReader call. io.py: 870 → 631 lines. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
…e_url SUPPORTED_IMAGE_EXTENSIONS was only used by the internal reader classes removed in the previous commit. FileReader.read_image_url() was never called in this branch — drop it and its now-unused base64 import. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
…ng in nvinference_client.py Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
… clean up OCRScoringQAStage by removing unused kwargs parameter. Signed-off-by: Ao Tang <aot@nvidia.com>
| data["image_path"] = self._get_image_path_str(task.data.image_path) | ||
| # Keep empty lists/strings/False (e.g. OCR may legitimately be []). | ||
| # Only drop fields that are explicitly None, and always omit is_valid. | ||
| self._file.write( | ||
| json.dumps({k: v for k, v in data.items() if v is not None and k != "is_valid"}, default=str) + "\n" | ||
| ) | ||
| elif self.valid_only: | ||
| self._skipped_count += 1 | ||
| return task | ||
| else: | ||
| data = task.data.to_dict() | ||
| data["image_path"] = self._get_image_path_str(task.data.image_path) | ||
| self._file.write( | ||
| json.dumps({k: v for k, v in data.items() if v is not None and k != "is_valid"}, default=str) + "\n" | ||
| ) | ||
| self._file.flush() # Flush after each write for safety | ||
| self._saved_count += 1 | ||
| return task |
There was a problem hiding this comment.
Invalid records silently lose
is_valid=False when written and re-read
ResultWriterStage.process always strips is_valid from the output via k != "is_valid". When valid_only=False (the CLI default — --valid-only is store_true defaulting to False), invalid records are written without an is_valid field. OCRData.from_dict defaults that field to True (data.get("is_valid", True)), so any code path that reloads the JSONL (e.g., a resume workflow or SkipProcessedStage) will treat every previously-invalid record as valid and re-attempt processing it.
The filter should preserve is_valid=False for invalid records while still omitting it for valid ones:
filtered = {k: v for k, v in data.items() if v is not None and k != "is_valid"}
if not task.data.is_valid:
filtered["is_valid"] = False
self._file.write(json.dumps(filtered, default=str) + "\n")Alternatively, include is_valid unconditionally and remove the k != "is_valid" filter.
| from nemo_curator.models.client import OpenAIClient | ||
|
|
||
|
|
||
| class NVInferenceModel(ModelInterface): |
There was a problem hiding this comment.
why is this a ModelInterface..
Can this just be a ProcessingStage?
To which we pass in a client which can be OpenAIClient or AsyncOpenAIClient
There was a problem hiding this comment.
This is just the modeling class. Similar to how we have in video for example that we have QwenLM class and then we invoke the model in CaptionGenerationStage.
Here we just have the NVInferenceModel which is supposed to work for any API endpoint from build.nvidia.com . and we use this model in the ModelProcessingStage (which is a ProcessingStage)
| logger.info(f"Tasks processed: {len(output_tasks)}") | ||
|
|
||
| merged = merge_output_shards(Path(args.output_path)) | ||
| logger.info(f"Output: {merged}") |
There was a problem hiding this comment.
We should probably add this to benchmark.. PS for that you'll also need to copy dataset for which this will work (so that we don't re download from the internet each time)..
copy the data to eos and dgx-a100 and then update the benchmarking/script/....yaml with the right path..
Feel free to ask @rlratzel for more questions around that
There was a problem hiding this comment.
I have the benchmark setup here in another branch: https://github.com/NVIDIA-NeMo/Curator/tree/aot/omni_sdg_benchmark
let's merge that later once we finalizing on the benchmark details so we know what to measure (i.e. GPU SKU)
…ncOpenAIClient. Introduced parameters for concurrent requests and updated related methods for async handling. Updated model processing stages to utilize ImageSampleTask instead of SingleDataTask for better image task management. Enhanced I/O stages to write JSONL results and read images from Hugging Face datasets. Improved OCR stages to handle image tasks more effectively. Updated tests to reflect changes in task handling and client initialization. Signed-off-by: Ao Tang <aot@nvidia.com>
Signed-off-by: Ao Tang <aot@nvidia.com>
|
/ok to test 7d4d113 |
… ModelInterface wrapper The OCR scoring path wrapped the NVIDIA Inference endpoint in NVInferenceModel(ModelInterface). ModelInterface is for GPU weight-bearing models (download weights, model_id_names = HF id); an HTTP client doesn't fit, and ~80% of the class re-implemented the AsyncOpenAIClient it already held — including a bespoke asyncio semaphore — while bypassing query_model and so losing the 429/connection retry+backoff that AsyncLLMClient provides for free. Replace it with a thin NVInferenceClient(AsyncOpenAIClient) that overrides only _query_model_impl (stream + reassemble delta.content for reasoning models, drop reasoning_content) and setup (resolve the API key from the env on the worker). It inherits the concurrency semaphore + retry. ModelProcessingStage now takes a client + model_name and calls client.query_model — mirroring the other SDG stages (e.g. QAMultilingualSyntheticStage) — with image→messages assembly moved into the stage. Nothing in the old wrapper was NVIDIA-endpoint-specific. Also: - drop the dead sync path (use_async was test-only) and the duplicated stream method - make setup() idempotent and forward stop/extra_kwargs in the create call - simplify build_prompt: ocr_dense is always OCRDenseItem, so drop the dict fallback - update unit tests for the client-based API Signed-off-by: Ao Tang <aot@nvidia.com>
| self._file.write( | ||
| json.dumps({k: v for k, v in data.items() if v is not None and k != "is_valid"}, default=str) + "\n" | ||
| ) |
There was a problem hiding this comment.
Invalid records written when
valid_only=False silently lose their is_valid=False state. OCRData.from_dict defaults is_valid to True (line 125 of ocr.py), so any downstream code that reloads the JSONL — e.g. a resume workflow or SkipProcessedStage — will treat every previously-invalid record as valid and reprocess it.
| self._file.write( | |
| json.dumps({k: v for k, v in data.items() if v is not None and k != "is_valid"}, default=str) + "\n" | |
| ) | |
| filtered = {k: v for k, v in data.items() if v is not None and k != "is_valid"} | |
| if not task.data.is_valid: | |
| filtered["is_valid"] = False | |
| self._file.write(json.dumps(filtered, default=str) + "\n") |
|
/ok to test 025bc07 |
|
/ok to test 85c82d5 |
Signed-off-by: Ao Tang <aot@nvidia.com>
|
/ok to test fc82dc0 |
|
/ok to test 146c38b |
Description
Adds the Nemotron OCR SDG pipeline — a multimodal synthetic data generation pipeline that converts images into structured OCR + QA conversation data for vision-language model training.
Pipeline stages
Key components
Tests
63 new unit tests across `tests/tasks/`, `tests/models/`, and `tests/stages/synthetic/omni/` — all CPU-only, no GPU required.
Checklist