# How Poly Data Ensures Pipeline Resumability: A Technical Deep Dive

> Poly Data ensures pipeline resumability by detecting last written records using line counts, cursor files, and timestamps. Restart your process seamlessly from the exact point of failure.

- Repository: [warproxxx/poly_data](https://github.com/warproxxx/poly_data)
- Tags: deep-dive
- Published: 2026-04-21

---

**Poly Data guarantees pipeline resumability by detecting the last successfully written record in each stage—using CSV line counts, persisted cursor files, and timestamp-hash filtering—and automatically continuing from that exact point on restart.**

The `warproxxx/poly_data` repository implements a robust three-stage data pipeline (market collection → Goldsky order-event scraping → live-trade processing) designed to survive interruptions without data loss or duplication. By leveraging idempotent write operations and explicit state persistence, the system ensures that rerunning any stage after a crash or manual stop simply continues from the previous checkpoint rather than starting over.

## Market Collection: Row Counting for API Offset

The first stage retrieves market data from an external API and appends it to `markets.csv`. To achieve **pipeline resumability** here, the code counts existing rows in the CSV and uses that count as the API pagination offset.

In [`update_utils/update_markets.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_markets.py), the helper `count_csv_lines` (lines 7-16) determines how many records have already been fetched:

```python
def count_csv_lines(filename):
    """Count lines in CSV to determine API offset."""
    if not os.path.exists(filename):
        return 0
    with open(filename, 'r') as f:
        return sum(1 for _ in f) - 1  # Subtract header

```

The main logic (lines 39-46) then passes this count as the `offset` parameter to the API request, ensuring the script fetches only new markets:

```python
current_offset = count_csv_lines('markets.csv')
print(f"Found {current_offset} existing records... Resuming...")

# Resume fetching from this offset

markets = api.get_markets(offset=current_offset, limit=batch_size)

```

This approach makes the market collection stage idempotent: running the script multiple times without new data on the API simply re-verifies the existing rows and exits gracefully.

## Goldsky Order-Event Scraping: Cursor Persistence with Sticky Timestamps

The second stage scrapes order-filled events from Goldsky and faces a complex resumability challenge: multiple events may share the same timestamp, requiring a "sticky" cursor that remembers not just the time but also the last processed ID within that timestamp window.

The solution resides in [`update_utils/update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_goldsky.py), which persists a cursor object (`timestamp`, `last_id`, `sticky_timestamp`) to [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json) after every successful batch.

### Cursor Retrieval Logic

The function `get_latest_cursor()` (lines 23-35) attempts to read the JSON state file first. If the file is missing, it falls back to parsing the last line of `orderFilled.csv` using `tail` (lines 81-86) to reconstruct the resume point:

```python
def get_latest_cursor():
    """Retrieve cursor from JSON or CSV fallback."""
    if os.path.exists('goldsky/cursor_state.json'):
        with open('goldsky/cursor_state.json', 'r') as f:
            return json.load(f)
    
    # Fallback: read last line of CSV

    last_line = os.popen('tail -n 1 orderFilled.csv').read().strip()
    if last_line:
        timestamp, order_id = parse_last_line(last_line)
        return {'timestamp': timestamp, 'last_id': order_id, 'sticky_timestamp': True}
    return None

```

### Cursor Persistence and Cleanup

After processing each batch, the cursor is atomically saved (lines 221-223), and old temporary files are cleaned up (lines 27-30):

```python

# After successful batch processing

with open('goldsky/cursor_state.json', 'w') as f:
    json.dump(cursor, f)

print(f"Resuming from cursor: {cursor['timestamp']}")

```

This "sticky" timestamp handling ensures that even when thousands of events share the same millisecond timestamp, the pipeline resumes from the exact event where it stopped, preventing gaps or duplicates in the event stream.

## Live-Trade Processing: Timestamp-Hash Deduplication

The third stage processes raw trades into cleaned `processed/trades.csv`. To maintain **pipeline resumability**, it uses the last written row to establish a watermark, filtering out any source data that has already been processed.

In [`update_utils/process_live.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/process_live.py), the detection logic (lines 11-14) reads the final line of the output CSV:

```python
if os.path.exists('processed/trades.csv'):
    # Get last line to determine resume point

    last_line = os.popen('tail -n 1 processed/trades.csv').read().strip()
    last_processed = parse_trade_line(last_line)
else:
    last_processed = None

```

The extraction of the last processed hash (lines 19-24) and the filtering logic (lines 47-60) then builds a deduplicated DataFrame:

```python

# Extract resume point details

last_timestamp = last_processed['timestamp']
last_hash = last_processed['trade_hash']

print(f"Resuming from: {last_timestamp}")

# Filter out already processed rows

df_process = df_source[
    (df_source['timestamp'] > last_timestamp) | 
    ((df_source['timestamp'] == last_timestamp) & 
     (df_source['trade_hash'] > last_hash))
]

```

This timestamp-hash combination acts as a composite primary key, ensuring that even if trades occur simultaneously, the pipeline distinguishes between processed and unprocessed records accurately.

## Summary

- **Market collection** uses `count_csv_lines` in [`update_utils/update_markets.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_markets.py) to set the API offset based on existing CSV rows, enabling pagination to resume exactly where it stopped.
- **Goldsky scraping** persists a cursor with sticky timestamps to [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json) via `get_latest_cursor()` and CSV fallbacks, implemented in [`update_utils/update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_goldsky.py).
- **Live-trade processing** reads the last row of `processed/trades.csv` using `tail` and filters source data by timestamp-hash combinations in [`update_utils/process_live.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/process_live.py).
- All three stages are idempotent, print explicit resume status messages, and handle rate-limiting with retries to ensure robust interruption recovery.

## Frequently Asked Questions

### What happens if the cursor JSON file is deleted in the Goldsky stage?

If [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json) is missing, [`update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_goldsky.py) automatically falls back to reading the last line of `orderFilled.csv` to reconstruct the resume point. This dual-layer persistence ensures that even without the JSON state file, the pipeline can determine where to resume by inspecting the actual data already written to disk.

### How does Poly Data handle duplicate timestamps when resuming?

The Goldsky implementation uses a "sticky timestamp" cursor that records both the `timestamp` and `last_id` within that timestamp window. In the live-trade stage, a composite key of timestamp and trade hash filters out duplicates. This ensures that events or trades sharing the same timestamp are processed exactly once, even when resuming mid-batch.

### Is the market collection stage safe to run multiple times?

Yes. The `update_markets` function counts existing rows in `markets.csv` and uses that count as the API offset. Running the script repeatedly without new data on the API simply results in zero new records being fetched and the existing file being rewritten identically, making the operation idempotent and safe for cron jobs or automated restarts.

### Can the pipeline resume if the CSV files are partially corrupted?

The code assumes standard CSV formatting for the last-line reads (using `tail`). If a file is corrupted such that the last line is incomplete, the parsing logic may fail. However, the Goldsky stage mitigates this by prioritizing the JSON cursor file over CSV inspection, providing a more reliable state source than raw CSV parsing alone.