Reading Data
Flowfile provides Polars-compatible readers with additional cloud storage integration and visual workflow features.
Polars Compatibility
All Flowfile readers accept the same parameters as Polars, plus optional description for visual documentation. The local CSV, Parquet, and Arrow IPC readers also take scan_mode and include_file_paths for reading a directory of files.
Local File Reading
CSV Files
import flowfile as ff
# Basic usage (same as Polars)
df = ff.read_csv("data.csv")
# With Flowfile description
df = ff.read_csv("data.csv", description="Load customer data")
# Polars parameters work identically
df = ff.read_csv(
"data.csv",
separator=",",
has_header=True,
skip_rows=1,
n_rows=1000,
description="Sample first 1000 customer records"
)
Key Parameters (same as Polars):
separator: Field delimiter (default:,)has_header: First row contains column names (default:True)skip_rows: Skip rows at start of filen_rows: Maximum rows to readencoding: File encoding (default:utf8)null_values: Values to treat as nullschema_overrides: Override column types
Parquet Files
# Basic usage
df = ff.read_parquet("data.parquet")
# With description
df = ff.read_parquet("sales_data.parquet", description="Q4 sales results")
Arrow IPC / Feather, NDJSON, and Avro Files
Flowfile also reads Arrow IPC/Feather, newline-delimited JSON (NDJSON), and Avro
files. These readers live in flowfile_frame and are not re-exported on the
ff namespace — import them directly:
from flowfile_frame import read_ipc, read_ipc_stream, read_ndjson, read_avro, scan_ipc, scan_ndjson
# Arrow IPC / Feather (lazy scan — like parquet)
df = read_ipc("data.arrow", description="Arrow IPC source")
# Arrow IPC stream (eager read — the footer-less format has no lazy scan)
df = read_ipc_stream("data.arrows")
# Newline-delimited JSON (lazy scan); a gzipped file is decompressed on the fly
df = read_ndjson("events.ndjson")
df = read_ndjson("events.ndjson.gz")
# Avro (eager read — offloaded to the worker so core never holds the dataset)
df = read_avro("data.avro")
IPC and NDJSON are scanned lazily, so they also provide scan_ipc / scan_ndjson.
Avro and the Arrow IPC stream format have no lazy scan in Polars, so their reads are offloaded
to the worker. read_csv and read_ndjson accept .gz files directly.
Reading a Directory of Files
read_csv, read_parquet, and read_ipc — and their scan_* aliases — read every matching file
in a directory as one table. Two parameters control this:
scan_mode:"single_file"or"directory". Auto-detected from the path ifNone(the default): an existing directory, a trailing separator, or a*,?, or[anywhere in the path selects"directory". HTTP(S) URLs are always read as a single fileinclude_file_paths: name of aStringcolumn holding each row's source file path.Noneor a blank name adds no column. Works in both scan modes; the name must not collide with a column already in the data
import flowfile as ff
# Every .csv under the folder, read as one frame and tagged with its source file
sales = ff.read_csv("data/monthly/", include_file_paths="source_file")
# An explicit glob is used as written
q1 = ff.read_csv("data/monthly/2026-0[123]-*.csv")
# Set the mode explicitly when the path does not exist yet
archive = ff.read_parquet("data/archive", scan_mode="directory")
sales stacks the rows of every file, with the absolute source path in the extra column:
| order_id | amount | source_file |
|---|---|---|
| 1 | 10.0 | /data/monthly/2026-01.csv |
| 2 | 20.0 | /data/monthly/2026-01.csv |
| 3 | 30.0 | /data/monthly/2026-02.csv |
A bare directory expands to a recursive glob for the format's extension — **/*.csv,
**/*.parquet, **/*.arrow, matched case-insensitively. An explicit pattern is used verbatim. A
pattern that matches no files raises NoFilesMatchedError.
Formats that support directory mode
CSV, Parquet, and Arrow IPC only. read_excel, read_ndjson, and read_avro do not take
scan_mode or include_file_paths. CSV additionally requires a UTF-8 encoding (utf8 or
utf8-lossy); any other encoding raises DirectoryScanUnsupportedError.
Parquet and Arrow IPC directory scans compare the column names and dtypes of every matched file before reading, and fail naming the file that disagrees. CSV uses Polars' own multi-file schema resolution.
The visual Read Data node exposes the same settings (Source →
Single file / Directory, plus File path column), and the cloud readers take the same scan_mode
parameter — see Cloud Storage Reading.
Scanning vs Reading
Flowfile provides both read_* and scan_* functions for Polars compatibility:
# These are identical in Flowfile
df1 = ff.read_csv("data.csv")
df2 = ff.scan_csv("data.csv") # Alias for read_csv
# The lazy IPC/NDJSON scans are imported from flowfile_frame (not ff.*)
from flowfile_frame import scan_ipc, scan_ndjson
df3 = scan_ipc("data.arrow")
df4 = scan_ndjson("events.ndjson")
Cloud Storage Reading
Flowfile extends Polars with specialized cloud storage readers that integrate with secure connection management.
Unified Cloud Storage Reader
read_from_cloud_storage() is a single entry point for reading any supported format from cloud storage. It dispatches to the appropriate format-specific reader internally.
import flowfile as ff
# Read Parquet (default format)
df = ff.read_from_cloud_storage(
"s3://bucket/data.parquet",
connection_name="my-conn",
)
# Read CSV
df = ff.read_from_cloud_storage(
"s3://bucket/data.csv",
file_format="csv",
connection_name="my-conn",
delimiter=",",
has_header=True,
)
# Read Delta with time travel
df = ff.read_from_cloud_storage(
"s3://warehouse/my_table",
file_format="delta",
connection_name="my-conn",
delta_version=5,
)
Parameters:
source: Cloud storage path (e.g.,s3://bucket/path/file.parquet)file_format:"csv","parquet","json", or"delta"(default:"parquet")connection_name: Name of the stored cloud storage connectionscan_mode:"single_file"or"directory". Auto-detected from path ifNonedelimiter: CSV field separator (default:;). Only used for CSVhas_header: Whether CSV has headers (default:True). Only used for CSVencoding: CSV encoding (default:utf8). Only used for CSVdelta_version: Delta table version for time-travel queries. Only used for Deltachanges_since: Read the table's change feed instead of its rows. Only used for Delta; see Delta Lake Readinginclude_change_preimage: Keep the before-image rows of each update in the change feed (default:False). Only used for Delta
Recommended Approach
read_from_cloud_storage() is the recommended way to read from cloud storage. The format-specific scan_* functions below still work and are useful when you want a more concise call for a known format.
Format-Specific Cloud Readers
Cloud CSV Reading
# Read from S3 with connection
df = ff.scan_csv_from_cloud_storage(
"s3://my-bucket/data.csv",
connection_name="my-aws-connection",
delimiter=",",
has_header=True,
encoding="utf8"
)
# Directory scanning (reads all CSV files)
df = ff.scan_csv_from_cloud_storage(
"s3://my-bucket/csv-files/",
connection_name="my-aws-connection"
)
CSV delimiter default
The cloud CSV readers default delimiter=";". Pass delimiter="," explicitly for comma-separated files (as above).
Cloud Parquet Reading
# Single file
df = ff.scan_parquet_from_cloud_storage(
"s3://data-lake/sales.parquet",
connection_name="data-lake-connection"
)
# Directory of files
df = ff.scan_parquet_from_cloud_storage(
"s3://data-lake/partitioned-data/",
connection_name="data-lake-connection",
scan_mode="directory"
)
Cloud JSON Reading
df = ff.scan_json_from_cloud_storage(
"s3://my-bucket/data.json",
connection_name="my-aws-connection"
)
Delta Lake Reading
# Latest version
df = ff.scan_delta(
"s3://data-lake/delta-table",
connection_name="data-lake-connection"
)
# Time-travel to a specific version
df = ff.scan_delta(
"s3://data-lake/delta-table",
connection_name="data-lake-connection",
version=5,
)
# Only the rows changed by commits after version 5
changes = ff.scan_delta(
"s3://data-lake/delta-table",
connection_name="data-lake-connection",
changes_since=5,
)
Parameters:
source: Cloud storage path of the Delta tableconnection_name: Name of the stored cloud storage connectionversion: Delta table version for time travelchanges_since: Read the table's change feed instead of its rows — anintreturns the changes committed after that version, an ISO-8601 string or adatetimethe changes committed at or after that instant. The result adds_change_type,_commit_versionand_commit_timestampinclude_change_preimage: Keep theupdate_preimagerows (the before-image of each update). Dropped by default
A change read needs the table to be change-tracked — write it with track_changes=True (see Writing Data) or turn tracking on from the visual reader. changes_since and version are mutually exclusive, changes_since="last_run" raises a ValueError because cursors exist for catalog tables only, and change reads are not supported for gs:// paths. Change Tracking covers how cloud paths differ from catalog tables and has a tested end-to-end example.
Catalog Reading
Read tables from the Flowfile catalog. The catalog provides a managed layer for discovering and versioning datasets stored as Delta tables. Both physical and virtual tables are supported.
Read a Table by Name
import flowfile as ff
# Read a catalog table (physical or virtual)
df = ff.read_catalog_table("my_table")
# Scope to a specific schema using a typed reference
schema = ff.CatalogReference("sales").schema("raw")
df = ff.read_catalog_table("my_table", schema=schema)
# Or call read_table directly on the schema handle
df = schema.read_table("my_table")
# Time travel to a specific Delta version (physical tables only)
df = ff.read_catalog_table("my_table", delta_version=5)
Parameters:
table_name: Name of the catalog table to read (required)schema: ASchemaReferenceidentifying the catalog/schema to read from. Preferred overnamespace_id.namespace_full_name: The schema as a plain"catalog.schema"string, resolved when the flow runs. Exported flow code uses this formdelta_version: Optional Delta version for time-travel queries (physical tables only)
Returns a FlowFrame. Use .collect() to materialize, .data to access the underlying LazyFrame, or open_graph_in_editor() to visualize in the UI.
Looking up tables by name
See Catalog References for the CatalogReference / SchemaReference API. The legacy namespace_id=<int> keyword is still accepted but discouraged — passing both schema= and namespace_id= raises ValueError.
Virtual table resolution
When reading a virtual table, the data is resolved on demand. Optimized virtual tables deserialize a stored execution plan instantly. Non-optimized virtual tables execute the producer flow to produce results. See Virtual Flow Tables for details.
Query with SQL
Use read_catalog_sql() to execute SQL queries against all catalog tables — both physical and virtual. Tables are registered by name in a Polars SQL context. read_catalog_sql is imported from flowfile_frame (it is not on the ff namespace):
from flowfile_frame import read_catalog_sql
# Query a single table
df = read_catalog_sql("SELECT * FROM customers WHERE region = 'Europe'")
# Join across catalog tables
df = read_catalog_sql("""
SELECT o.order_id, c.name, o.total
FROM orders o
JOIN customers c ON o.customer_id = c.id
WHERE o.total > 1000
""")
# Aggregate virtual and physical tables together
df = read_catalog_sql("""
SELECT category, SUM(amount) as total
FROM sales_summary
GROUP BY category
""")
Parameters:
sql_query: SQL query string to execute (required)
Returns a FlowFrame backed by a catalog SQL reader node. The SQL dialect is Polars SQL, which supports standard SELECT, WHERE, JOIN, GROUP BY, ORDER BY, HAVING, UNION, subqueries, and window functions.
Kafka Reading
Read messages from a Kafka topic using a stored Flowfile connection.
import flowfile as ff
df = ff.read_kafka(
"my-kafka-connection",
topic_name="events",
start_offset="earliest",
max_messages=10_000,
)
Parameters:
connection_name: Name of the stored Kafka connection (required)topic_name: Kafka topic to consume from (required)max_messages: Maximum number of messages to consume (default:100_000)start_offset: Where to start consuming:"earliest"or"latest"(default:"latest")poll_timeout_seconds: How long to poll for messages in seconds (default:30.0)value_format: Message value format (default:"json")
Returns a FlowFrame.
Database Reading
Read data from SQL databases using stored connections.
Setup Connection
import flowfile as ff
ff.create_database_connection(
connection_name="my_db",
database_type="postgresql",
host="localhost",
port=5432,
database="mydb",
username="user",
password="pass"
)
Read a Table
df = ff.read_database(
"my_db",
table_name="users",
schema_name="public"
)
Read with SQL Query
df = ff.read_database(
"my_db",
query="SELECT id, name FROM users WHERE active = true"
)
Parameters:
connection_name: Name of a stored database connection (required)table_name: Table to read fromschema_name: Database schema (e.g., "public")query: Custom SQL query (takes precedence overtable_name)
Return Type
read_database() returns a FlowFrame (not a raw Polars LazyFrame). The result supports .collect() to materialize data, .data to access the underlying LazyFrame, and open_graph_in_editor() to visualize the pipeline in the UI.
The tested integration example reads a table and a query from PostgreSQL through a stored connection:
big_budget = ff.read_database(
"analytics-postgres",
query="""
SELECT title, budget, vote_average
FROM public.movies
WHERE budget > 200000000
ORDER BY budget DESC
""",
).collect()
all_movies = ff.read_database(
"analytics-postgres", schema_name="public", table_name="movies"
).collect()
DuckDB
DuckDB connections point at a local file — no host, port, or credentials:
ff.create_database_connection(
connection_name="local_duckdb",
database_type="duckdb",
database="/path/to/analytics.duckdb",
)
df = ff.read_database("local_duckdb", table_name="events")
The tested DuckDB example (runs fully in-process):
high_rated = ff.read_database(
"local-duckdb",
query="SELECT title, rating FROM movies WHERE rating > 8.6 ORDER BY rating DESC",
).collect()
all_movies = ff.read_database("local-duckdb", table_name="movies").collect()
SQL Server
SQL Server connections use database_type="mssql" (default port 1433; the default schema is dbo):
ff.create_database_connection(
connection_name="analytics-sqlserver",
database_type="mssql",
host="sqlserver.example.com",
port=1433,
database="analytics",
username="user",
password="pass",
)
df = ff.read_database("analytics-sqlserver", schema_name="dbo", table_name="events")
The tested SQL Server example reads a table and a query through a stored connection:
top_rated = ff.read_database(
"analytics-sqlserver",
query="""
SELECT title, rating, director
FROM movies
WHERE rating >= 8.8
""",
).collect()
all_movies = ff.read_database(
"analytics-sqlserver", schema_name="dbo", table_name="movies"
).collect()
Denodo
Denodo connections use database_type="denodo" and go through Denodo's PostgreSQL-compatible
interface (default port 9996) with the standard psycopg2 driver. Use the Virtual DataPort virtual
database as database:
ff.create_database_connection(
connection_name="denodo-vdp",
database_type="denodo",
host="denodo.example.com",
port=9996,
database="analytics",
username="user",
password="pass",
ssl_enabled=True,
)
df = ff.read_database("denodo-vdp", query="SELECT customer_id, revenue FROM customer_sales")
Reads and queries are supported; writing to Denodo has not been verified against a live server.
Connection Management
Set up cloud and database connections once, then reference them by name. See Cloud Connection Management.
Examples
import flowfile as ff
# Read local file
customers = ff.read_csv("customers.csv", description="Customer master data")
# Read from cloud
orders = ff.scan_parquet_from_cloud_storage(
"s3://data-warehouse/orders/",
connection_name="warehouse",
description="Order history from data warehouse"
)
# Continue processing...
result = customers.join(orders, on="customer_id")