How Poly Data Ensures Pipeline Resumability: A Technical Deep Dive
Poly Data guarantees pipeline resumability by detecting the last successfully written record in each stage—using CSV line counts, persisted cursor files, and timestamp-hash filtering—and automatically continuing from that exact point on restart.
The warproxxx/poly_data repository implements a robust three-stage data pipeline (market collection → Goldsky order-event scraping → live-trade processing) designed to survive interruptions without data loss or duplication. By leveraging idempotent write operations and explicit state persistence, the system ensures that rerunning any stage after a crash or manual stop simply continues from the previous checkpoint rather than starting over.
Market Collection: Row Counting for API Offset
The first stage retrieves market data from an external API and appends it to markets.csv. To achieve pipeline resumability here, the code counts existing rows in the CSV and uses that count as the API pagination offset.
In update_utils/update_markets.py, the helper count_csv_lines (lines 7-16) determines how many records have already been fetched:
def count_csv_lines(filename):
"""Count lines in CSV to determine API offset."""
if not os.path.exists(filename):
return 0
with open(filename, 'r') as f:
return sum(1 for _ in f) - 1 # Subtract header
The main logic (lines 39-46) then passes this count as the offset parameter to the API request, ensuring the script fetches only new markets:
current_offset = count_csv_lines('markets.csv')
print(f"Found {current_offset} existing records... Resuming...")
# Resume fetching from this offset
markets = api.get_markets(offset=current_offset, limit=batch_size)
This approach makes the market collection stage idempotent: running the script multiple times without new data on the API simply re-verifies the existing rows and exits gracefully.
Goldsky Order-Event Scraping: Cursor Persistence with Sticky Timestamps
The second stage scrapes order-filled events from Goldsky and faces a complex resumability challenge: multiple events may share the same timestamp, requiring a "sticky" cursor that remembers not just the time but also the last processed ID within that timestamp window.
The solution resides in update_utils/update_goldsky.py, which persists a cursor object (timestamp, last_id, sticky_timestamp) to goldsky/cursor_state.json after every successful batch.
Cursor Retrieval Logic
The function get_latest_cursor() (lines 23-35) attempts to read the JSON state file first. If the file is missing, it falls back to parsing the last line of orderFilled.csv using tail (lines 81-86) to reconstruct the resume point:
def get_latest_cursor():
"""Retrieve cursor from JSON or CSV fallback."""
if os.path.exists('goldsky/cursor_state.json'):
with open('goldsky/cursor_state.json', 'r') as f:
return json.load(f)
# Fallback: read last line of CSV
last_line = os.popen('tail -n 1 orderFilled.csv').read().strip()
if last_line:
timestamp, order_id = parse_last_line(last_line)
return {'timestamp': timestamp, 'last_id': order_id, 'sticky_timestamp': True}
return None
Cursor Persistence and Cleanup
After processing each batch, the cursor is atomically saved (lines 221-223), and old temporary files are cleaned up (lines 27-30):
# After successful batch processing
with open('goldsky/cursor_state.json', 'w') as f:
json.dump(cursor, f)
print(f"Resuming from cursor: {cursor['timestamp']}")
This "sticky" timestamp handling ensures that even when thousands of events share the same millisecond timestamp, the pipeline resumes from the exact event where it stopped, preventing gaps or duplicates in the event stream.
Live-Trade Processing: Timestamp-Hash Deduplication
The third stage processes raw trades into cleaned processed/trades.csv. To maintain pipeline resumability, it uses the last written row to establish a watermark, filtering out any source data that has already been processed.
In update_utils/process_live.py, the detection logic (lines 11-14) reads the final line of the output CSV:
if os.path.exists('processed/trades.csv'):
# Get last line to determine resume point
last_line = os.popen('tail -n 1 processed/trades.csv').read().strip()
last_processed = parse_trade_line(last_line)
else:
last_processed = None
The extraction of the last processed hash (lines 19-24) and the filtering logic (lines 47-60) then builds a deduplicated DataFrame:
# Extract resume point details
last_timestamp = last_processed['timestamp']
last_hash = last_processed['trade_hash']
print(f"Resuming from: {last_timestamp}")
# Filter out already processed rows
df_process = df_source[
(df_source['timestamp'] > last_timestamp) |
((df_source['timestamp'] == last_timestamp) &
(df_source['trade_hash'] > last_hash))
]
This timestamp-hash combination acts as a composite primary key, ensuring that even if trades occur simultaneously, the pipeline distinguishes between processed and unprocessed records accurately.
Summary
- Market collection uses
count_csv_linesinupdate_utils/update_markets.pyto set the API offset based on existing CSV rows, enabling pagination to resume exactly where it stopped. - Goldsky scraping persists a cursor with sticky timestamps to
goldsky/cursor_state.jsonviaget_latest_cursor()and CSV fallbacks, implemented inupdate_utils/update_goldsky.py. - Live-trade processing reads the last row of
processed/trades.csvusingtailand filters source data by timestamp-hash combinations inupdate_utils/process_live.py. - All three stages are idempotent, print explicit resume status messages, and handle rate-limiting with retries to ensure robust interruption recovery.
Frequently Asked Questions
What happens if the cursor JSON file is deleted in the Goldsky stage?
If goldsky/cursor_state.json is missing, update_goldsky.py automatically falls back to reading the last line of orderFilled.csv to reconstruct the resume point. This dual-layer persistence ensures that even without the JSON state file, the pipeline can determine where to resume by inspecting the actual data already written to disk.
How does Poly Data handle duplicate timestamps when resuming?
The Goldsky implementation uses a "sticky timestamp" cursor that records both the timestamp and last_id within that timestamp window. In the live-trade stage, a composite key of timestamp and trade hash filters out duplicates. This ensures that events or trades sharing the same timestamp are processed exactly once, even when resuming mid-batch.
Is the market collection stage safe to run multiple times?
Yes. The update_markets function counts existing rows in markets.csv and uses that count as the API offset. Running the script repeatedly without new data on the API simply results in zero new records being fetched and the existing file being rewritten identically, making the operation idempotent and safe for cron jobs or automated restarts.
Can the pipeline resume if the CSV files are partially corrupted?
The code assumes standard CSV formatting for the last-line reads (using tail). If a file is corrupted such that the last line is incomplete, the parsing logic may fail. However, the Goldsky stage mitigates this by prioritizing the JSON cursor file over CSV inspection, providing a more reliable state source than raw CSV parsing alone.
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 →