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

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, 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.


# 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, 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.


# 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. 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, 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.


# 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 orchestrator executes them sequentially in the correct order:

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 fetches metadata to markets.csv, resuming via row count offsets.
  • Stage 2 (Event Scraping): update_goldsky.py pulls GraphQL events to goldsky/orderFilled.csv, using cursor_state.json for precise resumption.
  • Stage 3 (Trade Processing): process_live.py joins and enriches data into processed/trades.csv, skipping already-processed rows to ensure idempotency.
  • Orchestrator: 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. 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, 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.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →