How Poly Data Scrapes Order Events from Goldsky: Architecture and Implementation

Poly Data retrieves orderFilled events from Goldsky using a GraphQL-based scraper that implements sticky cursor pagination to handle high-throughput blockchain data, with support for both single-process and parallel segmented execution modes.

The warproxxx/poly_data repository provides a robust Python-based ingestion pipeline that synchronizes order book events from Goldsky's subgraph into local CSV storage. This article examines exactly how Poly Data scrapes order events from Goldsky, detailing the cursor-based pagination system, the dual-mode architecture, and the specific implementation details found in the source code.

Single-Process Scraping with update_goldsky.py

The single-process implementation provides a straightforward, sequential approach to event ingestion. Located in update_utils/update_goldsky.py, this module uses the scrape function to maintain a persistent loop that queries Goldsky, appends results to goldsky/orderFilled.csv, and manages state via goldsky/cursor_state.json.

The process initializes by calling get_latest_cursor, which inspects goldsky/cursor_state.json or the existing CSV tail to determine the last processed timestamp. The scrape function then enters a loop that continues until no new data remains, updating the cursor file after each successful batch write. When the scraper encounters the final batch (signaled by a response smaller than the requested batch size), it removes the cursor file to indicate completion.

from update_utils.update_goldsky import update_goldsky

# Execute the single-process scraper

update_goldsky()

Parallel Segmented Scraping with parallel_sync.py

For high-performance scenarios, parallel_sync.py implements a segmented parallel architecture that partitions the unsynced time range across multiple workers. The sync_segment function handles individual segment processing, while the main orchestration logic manages worker distribution and result merging.

The workflow begins with get_last_timestamp to determine the sync boundary, then splits the remaining range into equal segments based on the --workers parameter. Each worker executes sync_segment independently, using the same sticky pagination logic as the single-process version. After all workers complete, the merge_segments function concatenates temporary CSVs into the canonical goldsky/orderFilled.csv and updates the cursor state.


# Run with 5 concurrent workers (default)

python3 parallel_sync.py --workers 5

# Limit sync to a specific end timestamp for testing

python3 parallel_sync.py --workers 2 --end-ts 1767000000

Core Scraping Logic

Both scraping modes share identical core mechanics for interacting with the Goldsky GraphQL endpoint at https://api.goldsky.com/.../orderbook-subgraph/.../gn. The implementation follows a six-step pipeline that ensures reliable, resumable data ingestion.

Cursor Initialization and State Management

The scraper maintains persistence through goldsky/cursor_state.json, which stores three critical values: last_timestamp, last_id, and sticky_timestamp. The get_latest_cursor function (single-process) and get_last_timestamp function (parallel) read this state to establish the query starting point. This design allows the scraper to resume exactly where it left off after interruptions without re-fetching historical data.

When initializing, the system checks for existing cursor state. If found, it uses these values to construct the subsequent query's pagination parameters. The cursor updates atomically after each successful batch write via save_cursor, ensuring data durability even if the process terminates unexpectedly.

Building the GraphQL Where Clause (Sticky Pagination)

The query construction logic handles two distinct pagination modes. Standard mode uses timestamp_gt to fetch events newer than the last processed timestamp. Sticky mode activates when multiple events share the same Unix timestamp, using a compound filter of timestamp: X, id_gt: Y to paginate through events within that specific second.

This sticky cursor mechanism is essential for blockchain data where multiple transactions may occur in the same block timestamp. The following code demonstrates the where clause construction found in both update_utils/update_goldsky.py and parallel_sync.py:


# Standard pagination

where_clause = f'timestamp_gt: "{last_timestamp}"'

# Sticky pagination (when many events share the same timestamp)

where_clause = f'timestamp: "{sticky_timestamp}", id_gt: "{last_id}"'

query = f'''
{{
    orderFilledEvents(
        orderBy: timestamp, orderDirection: asc,
        first: {batch_size},
        where: {{{where_clause}}}
    ) {{
        id timestamp maker makerAmountFilled makerAssetId
        taker takerAmountFilled takerAssetId transactionHash
    }}
}}
'''

Query Execution and Error Handling

The goldsky_query function (in parallel_sync.py) and its equivalent in the single-process scraper execute the constructed GraphQL query against the Goldsky endpoint. The implementation includes exponential back-off retry logic to handle transient network failures or rate limiting, ensuring robust ingestion during high-load periods.

Each query requests a configurable batch_size of events, ordered by timestamp ascending. The system tracks whether the current timestamp has multiple events, setting the sticky_timestamp flag when necessary to trigger the compound pagination logic in subsequent iterations.

Data Processing and CSV Persistence

Upon receiving GraphQL responses, the scraper flattens the nested JSON structure, deduplicates records by id, and sorts the results by timestamp and id to maintain chronological order. The COLUMNS_TO_SAVE configuration determines which fields persist to disk.

Data writes append to goldsky/orderFilled.csv, with header management logic that omits the CSV header if the file already exists. In the parallel implementation, workers write to temporary segment files first, which merge_segments later concatenates into the main dataset while preserving sort order and removing duplicates across segment boundaries.

Summary

  • Dual-mode architecture: Poly Data offers both update_goldsky (single-process) and parallel_sync (multi-worker) implementations to accommodate different throughput requirements.
  • Sticky cursor pagination: The system uses compound filtering (timestamp + id_gt) to handle high-frequency events that share the same Unix timestamp, ensuring no data loss during ingestion.
  • Resumable state: Cursor persistence in goldsky/cursor_state.json enables fault-tolerant operation without duplicate data fetching.
  • GraphQL integration: Queries target the Goldsky orderbook subgraph, requesting specific event fields including maker/taker details and transaction hashes.
  • CSV output: Processed events append to goldsky/orderFilled.csv, with parallel modes utilizing temporary segment files that merge into the final dataset.

Frequently Asked Questions

What is the "sticky cursor" mechanism in Poly Data's Goldsky scraper?

The sticky cursor mechanism handles pagination when multiple orderFilled events share the same Unix timestamp. Instead of using only timestamp_gt, the scraper switches to a compound filter timestamp: X, id_gt: Y to iterate through events within that specific timestamp. This prevents data loss that would occur if the scraper simply advanced to the next second while leaving same-timestamp events unprocessed.

How does Poly Data resume scraping after an interruption?

The scraper persists state to goldsky/cursor_state.json after every successful batch write, storing the last_timestamp, last_id, and sticky_timestamp values. On startup, functions like get_latest_cursor or get_last_timestamp read this file to determine the query starting point. If the cursor indicates an active sticky state, the next query uses the compound where clause to continue from the exact event ID where processing stopped.

What is the difference between single-process and parallel scraping modes?

The single-process mode (update_utils/update_goldsky.py) processes events sequentially in a single loop, suitable for stable, low-latency environments. The parallel mode (parallel_sync.py) partitions the unsynced time range into segments, processes each concurrently via sync_segment workers, then merges results. Parallel mode significantly reduces sync time for large historical backfills but requires handling temporary CSV files and segment reconciliation.

How does the scraper handle GraphQL query failures?

Both implementations incorporate retry logic with exponential back-off when executing queries via goldsky_query or equivalent functions. If the Goldsky API returns errors or network timeouts, the scraper waits and retries rather than failing immediately. This ensures reliable ingestion during temporary service disruptions or rate limiting events.

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 →