awslabs / awslabs/graphrag-toolkit

[FEATURE] Document-level extraction manifest for distributed resume and incremental ingestion

Open
#325 0 comments 0 reactions 1 assignee Claimed by @noel-improv View on GitHub
enhancement
Dominant language
Python
Stars
442
Forks
106
Avg merge
2d 17h
Merged PRs (30d)
52

Description

> **⚠️ DRAFT — This issue is under discussion and not yet ready for implementation.**

## Problem Statement

The existing `Checkpoint`/`CheckpointFilter` tracks per-node completion via local filesystem savepoints:
```python
# CheckpointFilter.checkpoint_does_not_exist():
node_checkpoint_path = join(self.checkpoint_dir, tenant_node_id)
if os.path.exists(node_checkpoint_path): ...
```

Limitations at 1M+ document scale:

1. **Local filesystem only** — savepoints are files at `{output_dir}/save_points/{checkpoint_name}/{node_id}`. Doesn't work across distributed workers/containers.
2. **Node-level, not document-level** — tracks individual chunks. A document with 5 chunks = 5 savepoint files. Can't atomically answer "is this document fully extracted?"
3. **No failure state** — missing savepoint = "not done." No distinction between "never attempted" and "attempted, failed, needs retry."
4. **O(N) scan at resume** — at 1M docs × 5 chunks = 5M savepoint files. `os.path.exists()` per node at startup is slow.

### Why This Is Needed

When extraction or build fails mid-run (OOM, throttle, timeout, spot termination), the pipeline must re-process everything from scratch. At 1M+ documents that means re-calling Bedrock for documents already successfully extracted — burning time and money. The state store enables **resumability**: on restart, only incomplete documents are processed.

### What Already Exists

- `Checkpoint` + `CheckpointFilter` — per-node local savepoints ✅
- `S3BasedDocs` — S3 I/O for extraction output ✅
- `ExtractionPipeline.extraction_filters` — metadata-based filtering ✅
- `GraphStoreFactory.for_graph_store(uri)` — URI-based factory pattern ✅

## Proposed Solution: `DocumentStateStore`

A backend-agnostic abstraction that tracks per-document completion across both extraction and build phases:

```python
class DocumentStateStore(ABC):
"""Tracks per-document completion status across extraction/build phases."""

@abstractmethod
async def get_status(self, doc_id: str) -> Optional[DocStatus]:
...

@abstractmethod
async def set_status(self, doc_id: str, status: DocStatus, *, expected: Optional[DocStatus] = None) -> bool:
"""Conditional write. Returns False if expected != current (optimistic lock)."""
...

@abstractmethod
async def list_incomplete(self, phase: Phase) -> List[str]:
"""Return doc_ids not yet completed for the given phase."""
...
```

### Backend Implementations

| Backend | URI Scheme | Sweet Spot | Trade-offs |
|---------|-----------|------------|------------|
| **S3 JSONL** | `s3://bucket/prefix` | <100K docs | No extra infra. Full-file read on startup, append via PUT. ETag-based conditional writes. |
| **DynamoDB** | `dynamodb://table-name` | 1M–100M+ docs | O(1) per-doc read/write at any scale. ConditionExpression for optimistic locking. GSI on status for `list_incomplete`. |

### User Configuration

```python
GraphRAGConfig(
document_state_store='s3://my-bucket/state/' # S3 for <100K docs
# or
document_state_store='dynamodb://my-state-table' # DynamoDB for 1M+
)
```

Factory resolves from URI scheme — same pattern as existing `GraphStoreFactory.for_graph_store('neptune://...')`.

### Why Two Backends

Access pattern analysis:
- **S3 JSONL** — sequential reads (full manifest on startup), append-only writes. Works well when manifest fits in memory (<10MB ≈ 100K docs). No additional infrastructure.
- **DynamoDB** — O(1) random access per document. At 1M docs the JSONL file is ~100MB (multi-second startup reads), and S3 has no append API (read-modify-write races under concurrent writers). DynamoDB handles both problems natively.

Users pick based on their scale. The abstraction makes both first-class without forcing a runtime dependency on either.

### Operational Modes

- `resume` — process documents not in store + optionally retry `failed`
- `incremental` — new documents (not in store) only
- `retry-failed` — only documents with `status: failed`

### Integration

The store covers **both phases** — extraction and build. This eliminates the overlap between extraction-resume and build-resume (previously split across #323 and this issue). #323 now focuses solely on write pacing.

### Performance Impact

| Scenario | Current (node-level local) | Proposed (doc-level shared) |
|----------|----------------------------|------------------------------|
| Resume 1M corpus at 80% | Scan 5M savepoint files (minutes) | Query store for pending docs (ms per doc) |
| Add 1000 new docs to 1M corpus | Scan all savepoints to find new | Store diff: O(1000) |
| Distributed extraction (4 workers) | Not possible (local filesystem) | Shared store with conditional writes |
| Failure diagnosis | Scan all files | `list_incomplete(phase)` |

### Acceptance Criteria

- [ ] `DocumentStateStore` ABC with `get_status`, `set_status` (conditional), `list_incomplete`
- [ ] S3 JSONL implementation
- [ ] DynamoDB implementation
- [ ] URI-based factory (same pattern as `GraphStoreFactory`)
- [ ] Resume skips completed documents without re-scanning node checkpoints
- [ ] Incremental mode: only new documents processed
- [ ] Failed documents selectively retried
- [ ] Concurrent access safe for distributed workers (optimistic locking)
- [ ] Integrates with existing `Checkpoint` (can coexist for backward compat)
- [ ] Covers both extraction and build phases

## Alternatives Considered

- **Existing `Checkpoint`/`CheckpointFilter`** — works for single-machine, small-scale runs. Doesn't scale to distributed workers or 1M+ documents due to local filesystem dependency and O(N) scanning.
- **Querying the graph for already-built sources** — possible but adds load to Neptune during the resume check and couples extraction state to the graph store.
- **DynamoDB-only** — simpler single implementation but forces an infrastructure dependency on users with small corpora (<100K docs) who don't need it.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.