# Stages of the Poly Data Pipeline: A 3-Step ETL Guide

> Explore the 3 stages of the Poly Data pipeline: collect market data, scrape order events, and process trades. Understand this ETL guide for efficient data analysis.

- Repository: [warproxxx/poly_data](https://github.com/warproxxx/poly_data)
- Tags: how-to-guide
- Published: 2026-04-21

---

**The Poly Data pipeline consists of three resumable ETL stages: (1) Market Data Collection that fetches Polymarket metadata, (2) Order Event Scraping that pulls Goldsky GraphQL events, and (3) Trade Processing that enriches raw events into analysis-ready trades.**

The `warproxxx/poly_data` repository implements a continuous data ingestion system designed for Polymarket analysts. Understanding the stages of the Poly Data pipeline is critical for maintaining an up-to-date dataset without reprocessing historical data. Each stage operates independently with robust resumability, allowing the system to recover seamlessly from interruptions.

## Stage 1: Market Data Collection

The first stage ingests market metadata from the Polymarket REST API and persists it to a local CSV file. This foundation enables downstream enrichment by providing the context necessary to interpret raw trading events.

### Core Implementation

In [`update_utils/update_markets.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_markets.py), the `update_markets()` function fetches comprehensive market metadata including questions, outcomes, token contracts, timestamps, and cumulative volume. The implementation paginates through the API and appends new records to `markets.csv`.

```python

# Update markets (fetches new markets and appends to markets.csv)

from update_utils.update_markets import update_markets
update_markets()                     # optional: batch_size=500, csv_filename="markets.csv"

```

### Resumability Mechanism

Rather than re-fetching the entire catalog on every run, the stage counts existing rows in `markets.csv` to compute the correct `offset` parameter. This allows the pipeline to resume automatically from the last fetched record, ensuring idempotent incremental updates.

## Stage 2: Order Event Scraping

The second stage extracts raw trade events from the blockchain via Goldsky's GraphQL subgraph. This is the most data-intensive step, capturing every `orderFilledEvent` with millisecond granularity.

### GraphQL Pagination and Deduplication

Located in [`update_utils/update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_goldsky.py), the `update_goldsky()` function implements sophisticated pagination logic. It queries the subgraph by timestamp range and, when necessary, falls back to event ID pagination to handle high-frequency trading windows. The stage deduplicates events before appending them to `goldsky/orderFilled.csv`.

```python

# Scrape order-filled events from Goldsky

from update_utils.update_goldsky import update_goldsky
update_goldsky()                     # resumes automatically; writes to goldsky/orderFilled.csv

```

### Cursor State Management

Resumability relies on a cursor persisted in [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json). The system saves three critical values: `last_timestamp`, `last_id`, and `sticky_timestamp`. On restart, the loader reads this cursor (or falls back to parsing the last line of the CSV) and continues from the exact point where it stopped, preventing data gaps or duplicates.

## Stage 3: Trade Processing

The final stage transforms raw order events into enriched, analysis-ready trade records by joining blockchain data with the market metadata collected in Stage 1.

### Data Enrichment Logic

Implemented in [`update_utils/process_live.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/process_live.py), the `process_live()` function performs several critical transformations:

- **Asset Resolution:** Identifies the non-USDC asset in each trading pair
- **Market Mapping:** Resolves the specific market and outcome side using `markets.csv`
- **Direction Calculation:** Determines whether the taker is buying or selling
- **Normalization:** Computes standardized prices and normalized token amounts

The enriched rows are written to `processed/trades.csv`.

```python

# Process raw events into enriched trades

from update_utils.process_live import process_live
process_live()                       # appends new rows to processed/trades.csv

```

### Incremental Processing Guarantee

This stage ensures idempotency by examining the last processed row in `processed/trades.csv` (specifically the timestamp, transaction hash, maker, and taker). It skips all earlier rows, guaranteeing that re-running the pipeline never creates duplicate entries while always catching up from the last successful commit.

## Running the Complete Pipeline

While each stage can run independently, the [`update_all.py`](https://github.com/warproxxx/poly_data/blob/main/update_all.py) orchestrator executes them sequentially in the correct order:

```bash
uv run python update_all.py

```

This executes Market Data Collection, then Order Event Scraping, then Trade Processing, ensuring that `processed/trades.csv` reflects the most recent Polymarket activity.

## Summary

- **Stage 1 (Market Data):** [`update_markets.py`](https://github.com/warproxxx/poly_data/blob/main/update_markets.py) fetches metadata to `markets.csv`, resuming via row count offsets.
- **Stage 2 (Event Scraping):** [`update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_goldsky.py) pulls GraphQL events to `goldsky/orderFilled.csv`, using [`cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/cursor_state.json) for precise resumption.
- **Stage 3 (Trade Processing):** [`process_live.py`](https://github.com/warproxxx/poly_data/blob/main/process_live.py) joins and enriches data into `processed/trades.csv`, skipping already-processed rows to ensure idempotency.
- **Orchestrator:** [`update_all.py`](https://github.com/warproxxx/poly_data/blob/main/update_all.py) runs all three stages automatically.

## Frequently Asked Questions

### How does the Poly Data pipeline resume after an interruption?

Each stage implements a distinct resumability strategy. Stage 1 counts existing CSV rows in `markets.csv` to calculate the API offset. Stage 2 persists cursor state (`last_timestamp`, `last_id`, `sticky_timestamp`) to [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json). Stage 3 inspects the last row of `processed/trades.csv` to determine the resume point. These mechanisms ensure the pipeline restarts exactly where it left off without data loss or duplication.

### What is the difference between `orderFilled.csv` and `trades.csv`?

`goldsky/orderFilled.csv` contains raw blockchain events directly from the Goldsky subgraph, including transaction hashes, raw token amounts, and uncompressed market IDs. `processed/trades.csv` contains the final enriched output with resolved market questions, calculated directions (buy/sell), normalized prices, and human-readable asset symbols. The former is raw input; the latter is analysis-ready data.

### Where does the pipeline store intermediate state?

Stage 2 stores its cursor in [`goldsky/cursor_state.json`](https://github.com/warproxxx/poly_data/blob/main/goldsky/cursor_state.json), while Stage 1 and Stage 3 infer state directly from the output CSV files (`markets.csv` and `processed/trades.csv` respectively). This design minimizes external dependencies—no database is required—making the system portable and easy to inspect manually.

### Can I run individual stages without executing the full pipeline?

Yes. Each stage exposes a standalone entry point. Import `update_markets()` from `update_utils/update_markets`, `update_goldsky()` from `update_utils/update_goldsky`, or `process_live()` from `update_utils/process_live` to execute stages independently. This is useful when debugging specific data issues or when market metadata needs refreshing separately from event scraping.