Description
Describe the bug
When calling optimize(..., mode="append") with num_nodes > 1, pre-existing dataset chunk metadata is included in every node's intermediate {node_rank}-index.json file. When the master node merges these intermediate indices, the existing chunks are concatenated repeatedly, duplicating them N times in the final index.json, where N is the number of nodes.
Intended:
Existing [A, B]
├──► Node 0 [C] → 0-index.json ──┐
└──► Node 1 [D] → 1-index.json ──┴──► Merge → [A, B, C, D]
Current Bug:
Existing [A, B]
├──► Node 0 [A, B, C] → 0-index.json ──┐
└──► Node 1 [A, B, D] → 1-index.json ──┴──► Merge
↓
[A, B, C, A, B, D]
Root Cause
-
data_processor.py: Every node passes existing_index to _merge_no_wait().
-
writer.py: _merge_no_wait() adds existing_index["chunks"] to each node's local index.
-
data_processor.py: The master node concatenates all node index files during the final merge, causing the existing chunks to be repeated once per node.
Suggested Fix
The existing index should be excluded from intermediate node-level merges and added only once during the final global merge.
In DataProcessor._done():
existing_index = getattr(self, "existing_index", None)
merge_cache._merge_no_wait(
node_rank if num_nodes > 1 else None,
None if num_nodes > 1 else existing_index,
)
In DataProcessor._upload_index():
merge_cache._merge_no_wait(
existing_index=getattr(self, "existing_index", None),
)
Expected behavior
The final index.json should contain pre-existing chunks exactly once, followed by the newly generated chunks from all nodes:
rather than:
Description
Describe the bug
When calling
optimize(..., mode="append")withnum_nodes > 1, pre-existing dataset chunk metadata is included in every node's intermediate{node_rank}-index.jsonfile. When the master node merges these intermediate indices, the existing chunks are concatenated repeatedly, duplicating them N times in the finalindex.json, where N is the number of nodes.Root Cause
data_processor.py: Every node passesexisting_indexto_merge_no_wait().writer.py:_merge_no_wait()addsexisting_index["chunks"]to each node's local index.data_processor.py: The master node concatenates all node index files during the final merge, causing the existing chunks to be repeated once per node.Suggested Fix
The existing index should be excluded from intermediate node-level merges and added only once during the final global merge.
In
DataProcessor._done():In
DataProcessor._upload_index():Expected behavior
The final
index.jsonshould contain pre-existing chunks exactly once, followed by the newly generated chunks from all nodes:rather than: