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 (takerAmountFilled if taker is USDC, else makerAmountFilled)
  • token_amount: Selects the non-USDC side amount
  • price: Calculates as usd_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:

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:

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 →