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.pyfetches metadata tomarkets.csv, resuming via row count offsets. - Stage 2 (Event Scraping):
update_goldsky.pypulls GraphQL events togoldsky/orderFilled.csv, usingcursor_state.jsonfor precise resumption. - Stage 3 (Trade Processing):
process_live.pyjoins and enriches data intoprocessed/trades.csv, skipping already-processed rows to ensure idempotency. - Orchestrator:
update_all.pyruns 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →