How Poly Data Processes Raw Order Events Into Trades: A Complete Pipeline Breakdown
Poly Data converts low-level orderFilled events from the Goldsky GraphQL endpoint into clean, analysis-ready trade records using a deterministic, incremental ETL pipeline built with Python and Polars.
The warproxxx/poly_data repository implements a specialized ETL system that ingests raw Polymarket blockchain events and processes them into standardized trade formats. This pipeline transforms micro-unit raw data into human-readable CSVs while maintaining idempotency for daily incremental runs. Examining how Poly Data processes raw order events into trades reveals a sophisticated 10-stage transformation that handles asset resolution, market mapping, and price derivation.
Understanding the Raw Data Source
The pipeline consumes orderFilled events captured from the Goldsky subgraph and stored locally in goldsky/orderFilled.csv. These records represent low-level blockchain transactions where makers and takers exchange assets on Polymarket's CLOB (Central Limit Order Book).
The update_utils/update_goldsky.py script handles the upstream data ingestion, ensuring timestamp-ordered and deduplicated records before processing begins. Raw events contain Unix-epoch timestamps, micro-unit amounts (10⁶), and asset IDs that require significant transformation to become analysis-ready.
The 10-Stage Transformation Pipeline
Poly Data's core logic resides in update_utils/process_live.py, specifically within the get_processed_df function. The transformation executes through ten distinct stages:
1. Load Market Metadata
The pipeline begins by loading market definitions from markets.csv and missing_markets.csv via poly_utils/utils.py#get_markets. This creates a Polars DataFrame mapping token IDs to market IDs while identifying which token represents token1 versus token2 in each market pair.
2. Read and Parse Raw Events
The system streams goldsky/orderFilled.csv into a Polars DataFrame, converting Unix-epoch timestamps to proper datetime objects using pl.from_epoch(pl.col("timestamp"), time_unit="s"). This stage ensures temporal data is ready for time-series analysis.
3. Incremental Offset Processing
For idempotent daily runs, the pipeline checks for existing processed/trades.csv files. If found, it identifies the last processed row using a composite key of timestamp, hash, maker, and taker, then resumes processing only new events. This logic appears in process_live.py lines 11-30.
4. Identify the Non-USDC Asset
Each orderFilled event contains makerAssetId and takerAssetId fields. The pipeline determines which asset is the non-USDC token (the zero-value ID indicates USDC, Polymarket's quote asset) and creates a nonusdc_asset_id column to facilitate market mapping.
5. Join Market Information
The events table joins once with the long-format market DataFrame on nonusdc_asset_id. This join, executed in get_processed_df lines 35-42, yields the market_id and determines whether the non-USDC token represents the token1 or token2 side of the market.
6. Derive Human-Readable Asset Labels
The pipeline creates makerAsset and takerAsset columns, setting values to "USDC" when the corresponding ID is zero, otherwise inheriting the side label (token1 or token2). This transformation occurs in get_processed_df lines 44-48.
7. Normalize Micro-Unit Amounts
Raw amounts stored in micro-units (10⁶) require division by 10**6 to produce human-readable token quantities. This scaling applies uniformly across makerAmountFilled and takerAmountFilled columns.
8. Determine Trade Direction
Directional logic evaluates whether the taker is buying USDC. If takerAsset == "USDC", the trade is classified as BUY for the taker and SELL for the maker; otherwise, the inverse applies. This directional assignment appears in get_processed_df lines 58-75.
9. Compute USD, Token Amounts, and Price
Financial calculations derive three key metrics:
usd_amount: Selects the USDC side amount (takerAmountFilledif taker is USDC, elsemakerAmountFilled)token_amount: Selects the non-USDC side amountprice: Calculates asusd_amount / token_amount
This computation block occupies lines 77-95 in get_processed_df.
10. Assemble and Persist the Final Trade Table
The final stage selects standardized columns: timestamp, market_id, maker, taker, nonusdc_side, maker_direction, taker_direction, price, usd_amount, token_amount, and transactionHash. Results write incrementally to processed/trades.csv.
Implementation Example
Running the pipeline requires minimal code. First ensure the raw Goldsky data exists, then invoke the processing function:
import polars as pl
from update_utils.process_live import get_processed_df
# Load raw events with proper schema
schema = {"takerAssetId": pl.Utf8, "makerAssetId": pl.Utf8}
raw = pl.scan_csv("goldsky/orderFilled.csv", schema_overrides=schema).collect()
raw = raw.with_columns(
pl.from_epoch(pl.col("timestamp"), time_unit="s").alias("timestamp")
)
# Transform to trades
clean_trades = get_processed_df(raw)
print(clean_trades.head())
This produces a Polars DataFrame matching the output CSV structure, ready for pandas conversion or direct analysis.
Key Files and Functions
Understanding the repository structure clarifies how Poly Data processes raw order events into trades:
update_utils/process_live.py: Contains the main ETL orchestration andget_processed_dfimplementationpoly_utils/utils.py: Providesget_markets()for market metadata loading and deduplicationupdate_utils/update_goldsky.py: Handles Goldsky API ingestion intogoldsky/orderFilled.csvparallel_sync.py: Offers parallelized synchronization alternatives for high-volume scenarios
Summary
Poly Data's trade processing pipeline demonstrates a robust pattern for blockchain event transformation:
- Incremental ingestion resumes from the last processed row using composite keys, ensuring idempotent daily operations
- Asset resolution identifies non-USDC tokens and joins market metadata to map trades to specific markets
- Normalization converts micro-units to human-readable amounts while calculating directional indicators (BUY/SELL)
- Deterministic output produces standardized CSVs with price, USD amount, and token amount fields ready for quantitative analysis
Frequently Asked Questions
How does Poly Data handle duplicate order events?
The pipeline implements idempotent processing by tracking the last successful row using timestamp, transaction hash, maker, and taker fields. When resuming, it skips all previously processed rows, ensuring each orderFilled event results in exactly one trade record regardless of how many times the script runs.
What determines whether a trade is classified as BUY or SELL?
Direction classification depends on which party receives USDC. If the taker's asset is USDC (takerAsset == "USDC"), the taker is selling the non-USDC token, making it a SELL for the taker and a BUY for the maker. Conversely, if the maker receives USDC, the taker is buying, resulting in a BUY classification for the taker.
Why does the pipeline use Polars instead of Pandas?
Polars provides superior performance for the large CSV files generated by Goldsky scraping, particularly through its lazy evaluation (scan_csv) and efficient join operations. The pipeline leverages Polars' from_epoch function for timestamp conversion and its memory-efficient DataFrame operations for handling millions of order events.
Where does the raw orderFilled data originate?
The update_utils/update_goldsky.py script periodically queries the Goldsky GraphQL subgraph for orderFilledEvents, storing results in goldsky/orderFilled.csv. This file serves as the immutable input stream for the processing pipeline, ensuring the transformation remains reproducible and auditable.
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 →