Skip to content
This repository was archived by the owner on May 29, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/checks.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ jobs:
run: uv run pytest --cov --cov-report=xml

- name: Upload coverage
if: matrix.python-version == '3.12' && github.ref == 'refs/heads/main'
if: matrix.python-version == '3.12'
uses: codecov/codecov-action@v5
with:
token: ${{ secrets.CODECOV_TOKEN }}
12 changes: 12 additions & 0 deletions codecov.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
coverage:
status:
project:
default:
target: auto
patch:
default:
target: 80%

comment:
layout: "condensed_header, condensed_files, condensed_footer"
require_changes: true
4 changes: 4 additions & 0 deletions examples/cli/script.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,3 +7,7 @@
demo1 = DemoSource(name="demo1")
demo2 = DemoSource(name="demo2")
dag = il.DAG(demo1, demo2)

# Test with:
# uv run interloper run examples/cli/script.py --date 2026-01-01
# uv run interloper backfill examples/cli/script.py --start-date 2026-01-01 --end-date 2026-01-07
33 changes: 33 additions & 0 deletions examples/io/csv_.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
import datetime as dt
from pprint import pp
from typing import Any

import interloper as il

il.subscribe(print)

io = il.CsvIO(base_path="data")

partitioning = il.TimePartitionConfig(column="date")


@il.asset(io=io, partitioning=partitioning)
def a(context: il.ExecutionContext) -> list[dict[str, Any]]:
return [
{"date": context.partition_date, "value": 1},
{"date": context.partition_date, "value": 2},
]


@il.asset(io=io, partitioning=partitioning)
def b(context: il.ExecutionContext, a: list[dict[str, Any]]) -> list[dict[str, Any]]:
return [
{"date": context.partition_date, "value": int(a[0]["value"]) + 2},
{"date": context.partition_date, "value": int(a[1]["value"]) + 2},
]


dag = il.DAG(a, b)
dag.backfill(il.TimePartitionWindow(start=dt.date(2025, 1, 1), end=dt.date(2025, 1, 7)))

pp(b().partition_row_counts())
3 changes: 2 additions & 1 deletion packages/interloper/src/interloper/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
subscribe,
unsubscribe,
)
from interloper.io import IO, FileIO, IOContext, MemoryIO
from interloper.io import IO, CsvIO, FileIO, IOContext, MemoryIO
from interloper.normalizer import MaterializationStrategy, Normalizer
from interloper.partitioning import (
Partition,
Expand Down Expand Up @@ -82,6 +82,7 @@
"Config",
"ConfigError",
"ConfigSpec",
"CsvIO",
"DAGError",
"DAGSpec",
"DataNotFoundError",
Expand Down
2 changes: 2 additions & 0 deletions packages/interloper/src/interloper/io/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,14 @@
from interloper.io.adapter import DataAdapter, RowAdapter
from interloper.io.base import IO
from interloper.io.context import IOContext
from interloper.io.csv import CsvIO
from interloper.io.database import DatabaseIO, WriteDisposition
from interloper.io.file import FileIO
from interloper.io.memory import MemoryIO

__all__ = [
"IO",
"CsvIO",
"DataAdapter",
"DatabaseIO",
"FileIO",
Expand Down
140 changes: 140 additions & 0 deletions packages/interloper/src/interloper/io/csv.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
"""Filesystem-backed IO using CSV serialization."""

from __future__ import annotations

import csv
from pathlib import Path
from typing import Any

from interloper.io.base import IO
from interloper.io.context import IOContext
from interloper.partitioning.base import Partition, PartitionWindow
from interloper.serialization.io import IOSpec


class CsvIO(IO):
"""IO that reads and writes CSV files on the local filesystem.

Data is stored under ``{base_path}/{dataset}/{asset_name}/data.csv``
(or ``{base_path}/{asset_name}/data.csv`` when no dataset is set).
Partitioned assets add a ``{column}={id}`` subdirectory.

Data must be ``list[dict]`` — each dict represents a row, and the keys
of the first dict determine the CSV column headers.
"""

def __init__(self, base_path: str) -> None:
"""Initialize CsvIO.

Args:
base_path: Base directory path for CSV file storage.
"""
self.base_path = base_path

def _asset_path(self, context: IOContext) -> Path:
"""Return the base directory for an asset."""
return Path(self.base_path) / (context.asset.dataset or "") / context.asset.name

def _write_csv(self, file_path: Path, data: list[dict[str, Any]]) -> None:
"""Write a list of row dicts to a CSV file."""
file_path.parent.mkdir(parents=True, exist_ok=True)
if not data:
file_path.write_text("")
return
fieldnames = list(data[0].keys())
with file_path.open("w", newline="") as f:
writer = csv.DictWriter(f, fieldnames=fieldnames)
writer.writeheader()
writer.writerows(data)

def _read_csv(self, file_path: Path) -> list[dict[str, str]]:
"""Read a CSV file and return a list of row dicts.

Args:
file_path: Path to the CSV file.

Returns:
Rows as a list of dicts.

Raises:
FileNotFoundError: If the CSV file does not exist.
"""
if not file_path.exists():
raise FileNotFoundError(f"Data file not found: {file_path}")
with file_path.open(newline="") as f:
reader = csv.DictReader(f)
return list(reader)

def write(self, context: IOContext, data: list[dict[str, Any]]) -> None:
"""Write row data to a CSV file, creating partition subdirectories as needed."""
base = self._asset_path(context)

if context.partition_or_window is None:
self._write_csv(base / "data.csv", data)

elif isinstance(context.partition_or_window, PartitionWindow):
assert context.asset.partitioning
for partition in context.partition_or_window:
partition_path = base / f"{context.asset.partitioning.column}={partition.id}"
self._write_csv(partition_path / "data.csv", data)

else:
assert isinstance(context.partition_or_window, Partition)
assert context.asset.partitioning
partition_path = base / f"{context.asset.partitioning.column}={context.partition_or_window.id}"
self._write_csv(partition_path / "data.csv", data)

def read(self, context: IOContext) -> list[dict[str, str]] | list[list[dict[str, str]]]:
"""Read row data from a CSV file.

Returns:
A list of row dicts, or a list of lists for partition windows.
"""
base = self._asset_path(context)

if context.partition_or_window is None:
return self._read_csv(base / "data.csv")

elif isinstance(context.partition_or_window, PartitionWindow):
assert context.asset.partitioning
return [
self._read_csv(base / f"{context.asset.partitioning.column}={p.id}" / "data.csv")
for p in context.partition_or_window
]

else:
assert isinstance(context.partition_or_window, Partition)
assert context.asset.partitioning
return self._read_csv(
base / f"{context.asset.partitioning.column}={context.partition_or_window.id}" / "data.csv"
)

def partition_row_counts(self, context: IOContext) -> dict[str, int]:
"""Return row counts grouped by partition by scanning CSV files on disk."""
assert context.asset.partitioning is not None
column = context.asset.partitioning.column
base = self._asset_path(context)

counts: dict[str, int] = {}
if not base.exists():
return counts

for entry in sorted(base.iterdir()):
if entry.is_dir() and entry.name.startswith(f"{column}="):
partition_value = entry.name.split("=", 1)[1]
data_file = entry / "data.csv"
if data_file.exists():
rows = self._read_csv(data_file)
counts[partition_value] = len(rows)
return counts

def to_spec(self) -> IOSpec:
"""Convert to a serializable spec.

Returns:
The IOSpec representation of this CsvIO.
"""
return IOSpec(
path=self.path,
init={"base_path": self.base_path},
)
138 changes: 138 additions & 0 deletions packages/interloper/tests/io/test_csv.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,138 @@
"""Tests for CsvIO."""

import datetime as dt

import interloper as il


class TestCsvIO:
"""Tests for CsvIO."""

def test_initialization(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))
assert csv_io.base_path == str(tmp_path)

def test_write_read_non_partitioned(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset
def my_asset():
return [{"a": "1", "b": "2"}]

context = il.IOContext(asset=my_asset())
data = [{"a": "1", "b": "2"}, {"a": "3", "b": "4"}]

csv_io.write(context, data)
result = csv_io.read(context)
assert result == data

def test_write_read_partitioned(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset(partitioning=il.TimePartitionConfig(column="ds"))
def my_asset():
return [{"x": "1"}]

partition = il.TimePartition(dt.date(2025, 1, 1))
context = il.IOContext(asset=my_asset(), partition_or_window=partition)
data = [{"x": "1"}, {"x": "2"}]

csv_io.write(context, data)
result = csv_io.read(context)
assert result == data

def test_write_read_partition_window(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset(partitioning=il.TimePartitionConfig(column="ds"))
def my_asset():
return [{"v": "1"}]

window = il.TimePartitionWindow(
start=dt.date(2025, 1, 1),
end=dt.date(2025, 1, 3),
)
data = [{"v": "hello"}]

csv_io.write(il.IOContext(asset=my_asset(), partition_or_window=window), data)

# Reading a window returns a list per partition
result = csv_io.read(il.IOContext(asset=my_asset(), partition_or_window=window))
assert len(result) == 3
assert all(r == data for r in result)

def test_read_missing_file(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset
def my_asset():
return []

context = il.IOContext(asset=my_asset())
try:
csv_io.read(context)
assert False, "Expected FileNotFoundError"
except FileNotFoundError:
pass

def test_write_empty_data(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset
def my_asset():
return []

context = il.IOContext(asset=my_asset())
csv_io.write(context, [])
result = csv_io.read(context)
assert result == []

def test_partition_row_counts(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset(partitioning=il.TimePartitionConfig(column="ds"))
def my_asset():
return []

asset_instance = my_asset()

p1 = il.TimePartition(dt.date(2025, 1, 1))
p2 = il.TimePartition(dt.date(2025, 1, 2))

csv_io.write(il.IOContext(asset=asset_instance, partition_or_window=p1), [{"a": "1"}])
csv_io.write(il.IOContext(asset=asset_instance, partition_or_window=p2), [{"a": "1"}, {"a": "2"}, {"a": "3"}])

counts = csv_io.partition_row_counts(il.IOContext(asset=asset_instance))
assert counts == {"2025-01-01": 1, "2025-01-02": 3}

def test_partition_row_counts_empty(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset(partitioning=il.TimePartitionConfig(column="ds"))
def my_asset():
return []

counts = csv_io.partition_row_counts(il.IOContext(asset=my_asset()))
assert counts == {}

def test_to_spec(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))
spec = csv_io.to_spec()
assert spec.init == {"base_path": str(tmp_path)}
reconstructed = spec.reconstruct()
assert isinstance(reconstructed, il.CsvIO)
assert reconstructed.base_path == str(tmp_path)

def test_with_dataset(self, tmp_path):
csv_io = il.CsvIO(str(tmp_path))

@il.asset(dataset="my_dataset")
def my_asset():
return []

context = il.IOContext(asset=my_asset())
data = [{"col": "val"}]

csv_io.write(context, data)
result = csv_io.read(context)
assert result == data
Loading