# How Poly Data Processes Raw Order Events Into Trades: A Complete Pipeline Breakdown

> Discover how Poly Data transforms raw order events into trade records with its Python Polars ETL pipeline. Understand the complete processing breakdown for analysis-ready data.

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

---

**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`](https://github.com/warproxxx/poly_data/blob/main/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`](https://github.com/warproxxx/poly_data/blob/main/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`](https://github.com/warproxxx/poly_data/blob/main/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:

```python
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`](https://github.com/warproxxx/poly_data/blob/main/update_utils/process_live.py)**: Contains the main ETL orchestration and `get_processed_df` implementation
- **[`poly_utils/utils.py`](https://github.com/warproxxx/poly_data/blob/main/poly_utils/utils.py)**: Provides `get_markets()` for market metadata loading and deduplication
- **[`update_utils/update_goldsky.py`](https://github.com/warproxxx/poly_data/blob/main/update_utils/update_goldsky.py)**: Handles Goldsky API ingestion into `goldsky/orderFilled.csv`
- **[`parallel_sync.py`](https://github.com/warproxxx/poly_data/blob/main/parallel_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`](https://github.com/warproxxx/poly_data/blob/main/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.