Testing Pipelines with the Python SDK
This guide describes how to test Feldera SQL pipelines
using the Feldera Python SDK. This is an
introductory guide only; we try to keep the examples as simple
and readable as possible. The examples use plain Python functions
that check their results with assert and do not rely on a
Python unit testing framework.
A test builds a pipeline from a SQL program, sends input changes to the connectors of its tables, and compares the contents or the changes of its views with the expected results.
To run the examples, you need:
-
Feldera, running locally or remotely. The examples connect to its API at
http://localhost:8080. The quickest way to start a released version is the Docker image:docker run --pull always -p 8080:8080 --tty --rm -it images.feldera.com/feldera/pipeline-manager:latest -
The SDK:
pip install feldera.
Run Feldera from sources
To test against the current sources instead of a release, run Feldera
from a checkout of the
Feldera repository, without
Docker or Kubernetes. First install the prerequisites listed under
"Running Feldera from sources" in the repository's README.md. Then
build the SQL compiler, and start the pipeline manager from the root of
the repository:
(cd sql-to-dbsp-compiler && ./build.sh)
./scripts/start_manager.sh
The pipeline manager listens on http://localhost:8080 and runs the pipelines
on the same machine as the test. Thus a program can read inputs
from a local file with the file connector,
for example with "path": "file:///tmp/orders.json". With the Docker
image, the pipeline runs inside the container, so such a file must be
in a directory mounted into the container with -v.
Core concepts
The Python SDK gives access to the following abstractions:
| Concept | Meaning | Operations |
|---|---|---|
| Pipeline | A compiled SQL program together with its runtime. | Compile, start, start paused, pause, resume, stop, discard state. Read the status, the statistics, the errors, and the logs. |
| Table | An input of the program. | Insert rows, delete rows, replace the row with a given key. |
| View | An output of the program. | Follow its change stream. Read its snapshot if it is materialized. |
| Materialized table or view | A table or view that the pipeline stores in full. Tables with a PRIMARY KEY are always materialized. See Materialized tables and views. | Read its snapshot with an ad-hoc SQL query. |
| Snapshot | The current contents of a materialized table or view. | Read. |
| Connector | A link between a table or view and an external system. An input connector feeds a table; an output connector receives the data from a view. The SQL program declares most connectors. In addition, any HTTP client can send input to a table through its HTTP input connector, and can follow the changes of a table or view through an HTTP output connector. The SDK uses these connectors to push input and to listen. See Connectors. | Pause, start, read the statistics. |
| Change | A data row together with an integer weight: the number of copies inserted (positive) or deleted (negative). An update is a deletion of the old row and an insertion of the new row. The SDK represents a change of weight n as n changes of weight +1 or -1, each with an insert_delete column that holds the weight. | Send to a table as input. Receive from a change stream as output. |
| Step | The pipeline processes a batch of input changes and computes the resulting changes of all views during a step. After each step, every view agrees with all input received so far. The pipeline can process one input in several steps. | Wait until all steps that process an input are complete. |
| Transaction | A group of input changes that the pipeline processes as one unit. All input changes that the pipeline receives between the start and the commit belong to the transaction. The views change once, at commit. See Transactions. | Start, commit, read the status. |
| Change stream | The changes of a table or view, in the order of the steps that produce them. | Connect a listener, read the changes received so far. |
Lifecycle of a test
A test takes a pipeline through these phases:
| Phase | Test does | SDK calls | Pipeline state |
|---|---|---|---|
| 1 | Compile the program | PipelineBuilder(client, name, sql).create_or_replace() returns a Pipeline object | Stopped |
| 2 | Start the pipeline paused, connect listeners | pipeline.start_paused(), pipeline.listen(view) | Paused |
| 3 | Resume the pipeline | pipeline.resume() | Running |
| 4 | Push input, then wait until the pipeline processes it | pipeline.input_json(table, changes) | Running |
| 5 | Read snapshots and changes | pipeline.query(sql), listener.to_dict() | Running |
| 6 | Stop the pipeline, discard its state | pipeline.stop(force=True), pipeline.clear_storage() | Stopped |
| 7 | Delete the pipeline | pipeline.delete(clear_storage=True) | Deleted |
A compiled pipeline can go through phases 2 to 6 many times, so many tests can use one compilation.
Compilation is usually the slowest phase, and Feldera compiler server can only compile a limited number of programs at the same time. The compiler server may be the bottleneck of the test suite. Compile each program once and try to reuse a pipeline in multiple tests. Several tests can also be implemented as a single big SQL program.
A test that reads only snapshots can combine phases 2 and 3: it can
start the pipeline directly with pipeline.start(), as in
A first test. stop(force=True) stops the pipeline
at once, without making a checkpoint. clear_storage() discards the
state of the pipeline, enabling the next test to start with empty tables.
We recommend deleting a pipeline only after tests pass. Keeping a pipeline after a failure allows you to inspect the program, status, and errors (see Common problems).
A first test
This test is a Python program that does the following:
- It compiles a SQL program that computes the total of the orders of each customer.
- It starts a pipeline that runs the SQL program.
- It inserts three orders.
- It reads the totals and compares them with the expected values.
- It stops the pipeline and deletes it.
Save the program in a file called test_orders.py:
from feldera import FelderaClient, PipelineBuilder
SQL = """
CREATE TABLE orders (
id BIGINT NOT NULL PRIMARY KEY,
customer VARCHAR NOT NULL,
amount INT NOT NULL
);
-- MATERIALIZED: the pipeline stores the entire view, so the test
-- can read all of its contents at any time with an ad-hoc query.
CREATE MATERIALIZED VIEW customer_totals AS
SELECT customer, SUM(amount) AS total, COUNT(*) AS num_orders
FROM orders
GROUP BY customer;
"""
def check_totals_per_customer() -> None:
client = FelderaClient("http://localhost:8080")
pipeline = PipelineBuilder(client, name="test-orders", sql=SQL).create_or_replace()
pipeline.start()
try:
# Blocks until the pipeline has processed the rows and the results are visible.
pipeline.input_json(
"orders",
[
{"id": 1, "customer": "alice", "amount": 10},
{"id": 2, "customer": "bob", "amount": 5},
{"id": 3, "customer": "alice", "amount": 7},
],
)
rows = list(pipeline.query("SELECT * FROM customer_totals ORDER BY customer"))
assert rows == [
{"customer": "alice", "total": 17, "num_orders": 2},
{"customer": "bob", "total": 5, "num_orders": 1},
]
finally:
pipeline.stop(force=True)
# Only reached when the check passed: a failed test keeps the stopped
# pipeline for inspection.
pipeline.delete(clear_storage=True)
if __name__ == "__main__":
check_totals_per_customer()
You can run this with python test_orders.py. If the check fails, the program
stops with an AssertionError. The program goes through the following phases:
| Phase | Code | Result |
|---|---|---|
| 1 | PipelineBuilder(...).create_or_replace() | Feldera compiles the program. The first compilation can be slow, but later compilations are much faster. |
| 2 | pipeline.start() | The pipeline starts with empty tables. |
| 3 | pipeline.input_json("orders", [...]) | Insert three rows into orders and update customer_totals. The pipeline.input_json() call returns only after the pipeline has processed the rows and sent the resulting changes to all its outputs, so the new contents of customer_totals are visible. |
| 4 | pipeline.query("SELECT ...") | Read the current snapshot of customer_totals and compare it with the expected rows. |
| 5 | pipeline.stop(force=True) | The pipeline stops. The finally block makes sure that this also occurs when the check fails. |
| 6 | pipeline.delete(clear_storage=True) | Feldera discards the state of the pipeline and deletes it. This runs only when the check passed: after a failure, the stopped pipeline remains, and you can inspect its program, status, and errors in the Web Console or with the SDK (see Common problems). |
You can query the data in a materialized view with
ad-hoc queries. (Ad-hoc queries use a different SQL
dialect than Feldera SQL programs, so we recommend keeping them simple,
for example SELECT * FROM customer_totals). Ad-hoc queries also
use the same computational resources as the pipeline, so running
expensive ad-hoc queries on large data can interfere or even
crash the pipeline.
Because input_json() returns only after
the pipeline has processed the rows, the new contents of
customer_totals are visible as soon as the call completes.
Thus the ad-hoc query in phase 4 sees the totals of all three orders.
Send and receive data
A test sends data to the input connectors of the tables, and receives data from the tables and views. Rows are Python dictionaries whose keys are the column names. These examples use the program from A first test:
# A row of the 'orders' table:
{"id": 1, "customer": "alice", "amount": 10}
# A list of input changes using the "insert_delete" format. These changes delete a row and
# insert another one: together, the two changes update order 1.
[
{"delete": {"id": 1, "customer": "alice", "amount": 10}},
{"insert": {"id": 1, "customer": "alice", "amount": 12}},
]
# A change of the 'customer_totals' view, as returned by a listener:
{"customer": "alice", "total": 10, "num_orders": 1, "insert_delete": 1}
These methods can be used to send data to a pipeline:
| Method | Sends | Notes |
|---|---|---|
pipeline.input_json(table, changes) | A list of changes | When using update_format="insert_delete", each element is {"insert": row} or {"delete": row}. See Updates and deletes. |
pipeline.input_pandas(table, df) | The rows of a pandas DataFrame | Inserts the rows (delete is unsupported). |
pipeline.execute("INSERT INTO ...", wait=True) | The rows of an ad-hoc INSERT statement | Inserts only: ad-hoc queries cannot delete rows. To delete rows, use pipeline.input_json(). |
| A connector declared in SQL | Data from an external system, for example a file or a Kafka topic | The test writes the data to the external system, which sends it to the pipeline. |
| The datagen connector, declared in SQL | Rows that the connector generates from a plan in its configuration | See Tests with fixed inputs. |
These methods receive data:
| Method | Receives | Notes |
|---|---|---|
pipeline.query("SELECT ...") | The snapshot: the current contents of the queried tables and views | Only for materialized tables and views. |
pipeline.listen(view) | A listener: an OutputHandler object that receives the change stream of one table or view | For all tables and views. See the listener methods below. |
| A connector that the program declares | Changes that the pipeline writes to an external system | The test reads from the external system. |
pipeline.listen() returns a listener. The listener is essentially a queue,
which automatically receives changes from the monitored table or view.
The listener interacts with the pipeline using an HTTP output connector,
which is automatically attached to every table and view.
The listener has two methods that return the changes that it received so far:
| Listener method | Returns |
|---|---|
listener.to_dict() | The changes as a list of dictionaries. Each change contains the columns of its row and an additional insert_delete column, with the value +1 or -1. |
listener.to_pandas() | The same changes as a pandas DataFrame. |
Both methods return only the changes that have been received since the previous call, and remove them from the listener's queue. These two calls return immediately.
For example, when a second order of alice updates
her total, listener.to_pandas() on customer_totals returns a
deletion of the old row and an insertion of the new row:
customer total num_orders insert_delete
0 alice 10 1 -1
1 alice 17 2 1
listener.to_dict() returns the same changes encoded as Python dictionaries.
A change of weight n appears as n identical changes of weight +1 or
-1. When no change has arrived,
listener.to_pandas() returns an empty DataFrame without columns.
The listener returns every change that it receives, in the order of
arrival, and does not combine multiple changes.
For example, if one step inserts a row and
a later step deletes it, listener.to_dict() returns both the insertion
and the deletion. Within one step, however, the pipeline itself combines
changes, so the changes within one step never contain
both an insertion and a deletion of the same row.
A listener uses a background thread to read data from the pipeline.
The pipeline waits for this thread, so a thread that cannot read fast
enough slows the pipeline down. The thread stores every change in the
listener's queue, which has no size limit, so the Python test may run
out of memory when the pipeline produces changes faster than the
listener can absorb them.
The next section explains when the data that a test receives reflects a particular set of input changes.
Synchronization
A pipeline runs in its own process, independently of the test: the input connectors feed its tables, it processes the input in steps, and it sends the results to its output connectors. The test does not write to the tables directly. It writes to input connectors, either through the SDK or through an external system that a connector reads.
The computation model of the pipeline is strongly consistent. The pipeline processes its input in steps, and after each step the contents of every view agree with all the input that the pipeline has received up to that step. See A synchronous streaming model.
However, this guarantee applies to the state of the pipeline, and does not necessarily apply to the interaction of the test with the pipeline. The test observes the pipeline through separate operations, and each operation can observe the state of the pipeline at a different moment:
| When running | Observe |
|---|---|
| One ad-hoc query | The state after one step, for all the tables and views that the query reads. |
| Two ad-hoc queries | The two queries may observe the same or different steps. |
| A listener | The changes of one table or view, in the order of the steps that produce them, delivered after the listener has connected. The changes of one step can arrive in several chunks; a chunk never contains changes of two steps. |
| Independent listeners | Separate change streams. |
Thus the data that a test gets from the pipeline reflects a consistent snapshot only when the test synchronizes with the pipeline. Synchronization makes sure that:
- The observed state includes all the input that the test sent, and no input that the test has not sent yet.
- Separate reads observe the same state. After the pipeline has processed all its input, and when no connector delivers new input, the contents of the tables and views do not change.
Without proper synchronization, tests may be flaky.
For this reason, a test usually sends all the input itself, so that it
controls when input arrives.
A test can control the values of the real-time clock as described in
Testing programs that use the real-time clock NOW().
Feldera pipelines are deterministic: given the same input, a pipeline
produces the same output. Thus a
test that sends the same input always gets the same result, which makes
tests repeatable. However, the ad-hoc queries used to query tables or
views are not necessarily deterministic (e.g., a LIMIT query).
We recommend against using elapsed time for synchronization: a test that sleeps for a fixed time and then reads fails when the pipeline is slower than expected. Tests should wait for a deterministic condition instead. Some conditions are enforced by a blocking SDK call. For other conditions, the test polls the condition, as the helper functions do. Timeouts are used to deal with unresponsive pipelines or unresponsive systems under test (e.g., Kubernetes does not allocate enough resources for the pipeline to start):
| To read | Method | Example |
|---|---|---|
| A snapshot | Call a method that waits until the pipeline has processed the input, and query when it returns. See the table below. | A first test |
| Change stream | Connect the listeners before the first step of the pipeline, then wait for a known number of changes. | Check the contents of a view as a change stream |
| Absence of a change | Push a second input with a known effect, and wait for its change. | Check that a change does not occur |
Result of NOW() | Wait until a view that shows NOW() has the new value. | Testing programs that use the real-time clock NOW() |
The following calls return only after the pipeline has processed the input:
| Call | Returns when | After it returns |
|---|---|---|
pipeline.input_json(table, changes), pipeline.input_pandas(table, df) | The pipeline processed the input and sent the resulting changes to all its outputs. | pipeline.query() shows the effect of the input. |
pipeline.execute("INSERT ...", wait=True) | Same as pipeline.input_json(). | Same as pipeline.input_json(). |
pipeline.commit_transaction() | The commit is complete. | pipeline.query() shows the effect of the transaction. |
pipeline.wait_for_completion() | Every input connector signaled the end of its input, and the pipeline processed all of it. See below. | pipeline.query() shows the effect of all input. |
If you want to test programs without materialized views:
- Connect the listeners while the pipeline is paused, before it
processes any input: call
start_paused(), thenlisten(), thenresume(). This ensures that listeners receive all changes, including from the first step. - Wait for the changes. Beware: the pipeline sends each change to the
listener over the network, and a background thread of the test
receives it. Thus the pipeline can send a change, and
input_json()can return, before the corresponding output change reaches the listener.to_dict()returns only the changes already in the queue. Theread_changes()function in Helper functions waits until a specified number of changes arrive.
End of input
Some input connectors have a notion of "end of input", for example
a file connector or the datagen connector. A
connector whose source can produce an unbounded data stream, for example a
Kafka topic, may never reach an "end of input". Some
connectors support both bounded and unbounded stream modes: a file connector
with configuration follow: true does not emit an end of input
(see File input connector configuration).
A Delta Lake connector signals "end of input" only after it reads a snapshot
(mode: snapshot), or after it reaches the specified end_version in the
connector configuration while following the transaction log
(see Delta Lake input connector configuration).
The HTTP connectors do not signal "end of input".
pipeline.wait_for_completion() returns when every input connector
except the HTTP connectors has signaled "end of input", and the pipeline
has processed all its input and sent the resulting changes to all its
outputs. It never returns if a connector does not signal "end of
input", or if the program uses NOW() and the clock follows the system
clock.
Helper functions
The following functions are not part of the SDK, but may be handy.
wait_for_condition() polls a condition specified as a Python predicate
function, read_changes() waits until a listener has received a number
of changes, and sorted_dicts() sorts rows or changes.
import time
from collections.abc import Callable, Iterable, Mapping
from typing import Any
from feldera.output_handler import OutputHandler
def wait_for_condition(
description: str,
predicate_func: Callable[[], bool],
timeout_s: float | None = None,
poll_interval_s: float = 0.1,
) -> None:
"""Keep re-evaluating `predicate_func` until it returns True or the timeout elapses.
:param description: Human-readable description used in the timeout error.
:param predicate_func: Callable returning True when a condition is met.
:param timeout_s: Maximum wait time in seconds. None means wait forever.
:param poll_interval_s: How frequently to call the predicate, in seconds.
:raises TimeoutError: If the condition is not met within `timeout_s`.
"""
timestamp_deadline_s = (
time.monotonic() + timeout_s if timeout_s is not None else float("inf")
)
while not predicate_func():
if time.monotonic() > timestamp_deadline_s:
raise TimeoutError(
f"timeout ({timeout_s:.1f}s) waiting for condition '{description}'"
)
time.sleep(poll_interval_s)
def read_changes(
listener: OutputHandler, count: int, timeout_s: float | None = None
) -> list[dict[str, Any]]:
"""Wait until `listener` has received at least `count` changes; return all
received changes (could be more than count).
Each change has an `insert_delete` value of +1 or -1; a change of weight n
arrives as n changes.
:param listener: The listener that `Pipeline.listen()` returned.
:param count: Number of changes to wait for.
:param timeout_s: Maximum wait time in seconds. None means wait forever.
:raises TimeoutError: If fewer than `count` changes arrive within `timeout_s`.
"""
wait_for_condition(
f"{count} change(s) on the '{listener.view_name}' listener",
lambda: len(listener.to_pandas(clear_buffer=False)) >= count,
timeout_s,
)
return listener.to_dict()
def sorted_dicts(items: Iterable[Mapping[str, Any]]) -> list[Mapping[str, Any]]:
"""Return `items`, a list of dictionaries such as rows or changes, sorted
using their natural order."""
return sorted(items, key=lambda item: sorted(item.items()))
Test scenarios
Here are solutions to some common test scenarios.
Each scenario is implemented in a Python function that receives a
compiled pipeline. It starts the pipeline, and at the end it stops
the pipeline and discards its state, preparing the pipeline for a
subsequent test. Unless a scenario shows its own program, it uses the
SQL program from A first test.
The following program creates the pipeline and invokes two scenarios. It deletes the pipeline only when every scenario passed:
from feldera import FelderaClient, Pipeline, PipelineBuilder
client = FelderaClient("http://localhost:8080")
pipeline = PipelineBuilder(client, name="test-orders", sql=SQL).create_or_replace()
try:
check_totals_in_any_order(pipeline)
check_second_order_updates_total(pipeline)
finally:
pipeline.stop(force=True)
# Only reached when every scenario passed.
pipeline.delete(clear_storage=True)
Check the expected contents of a view
Read the snapshot of a view using the query() method after the input_json() call returns.
query() returns the rows in no defined order. To compare them with an expected result, you can:
- Use an ad-hoc query that sorts the rows using
ORDER BY, and compare them with a list in the same order, as A first test does. - Sort the rows with
sorted_dicts()above.
def check_totals_in_any_order(pipeline: Pipeline) -> None:
pipeline.start()
pipeline.input_json(
"orders",
[
{"id": 1, "customer": "alice", "amount": 10},
{"id": 2, "customer": "bob", "amount": 5},
],
)
rows = list(pipeline.query("SELECT * FROM customer_totals"))
pipeline.stop(force=True)
pipeline.clear_storage()
assert sorted_dicts(rows) == sorted_dicts(
[
{"customer": "alice", "total": 10, "num_orders": 1},
{"customer": "bob", "total": 5, "num_orders": 1},
]
)
Note: query() can only be used to read the contents of materialized
tables and views.
Check the contents of a view as a change stream
To inspect a view that is not materialized, you need to capture all the changes that it produces.
Start the pipeline paused, connect the listener, and then resume the
pipeline, so that the listener connects before the first step. Then
wait for the changes with read_changes().
def check_second_order_updates_total(pipeline: Pipeline) -> None:
# Connect the listener before the first step.
pipeline.start_paused()
totals = pipeline.listen("customer_totals")
pipeline.resume()
pipeline.input_json("orders", [{"id": 1, "customer": "alice", "amount": 10}])
pipeline.input_json("orders", [{"id": 2, "customer": "alice", "amount": 7}])
changes = read_changes(totals, 3)
pipeline.stop(force=True)
pipeline.clear_storage()
assert sorted_dicts(changes) == sorted_dicts(
[
{"customer": "alice", "total": 10, "num_orders": 1, "insert_delete": 1},
{"customer": "alice", "total": 10, "num_orders": 1, "insert_delete": -1},
{"customer": "alice", "total": 17, "num_orders": 2, "insert_delete": 1},
]
)
The pipeline processes each input_json() call in one or more steps
of its own; the next call starts only after the pipeline has processed
the previous one. The first call inserts the row for alice. The second call
updates the row: it deletes the old row and inserts the new row.
read_changes() returns all changes that the listener received, which
can be more than count.
Updates and deletes
The SDK can express a change using the following formats:
| Change | update_format | Sample element of the data list |
|---|---|---|
| Insert | "raw" (default) | {"id": 1, "customer": "alice", "amount": 10} |
| Insert | "insert_delete" | {"insert": {"id": 1, ...}} |
| Delete | "insert_delete" | {"delete": {"id": 1, ...}} |
| Replace a row (table with a primary key) | any | An insert with the key of an existing row. |
| Change some columns (table with a primary key) | "insert_delete" | {"update": {"id": 1, "amount": 12}}. See the JSON format. |
def check_replace_and_delete(pipeline: Pipeline) -> None:
pipeline.start()
pipeline.input_json(
"orders",
[
{"id": 1, "customer": "alice", "amount": 10},
{"id": 2, "customer": "bob", "amount": 5},
],
)
# This insert has the key of order 1, so it replaces order 1.
pipeline.input_json("orders", [{"id": 1, "customer": "alice", "amount": 25}])
pipeline.input_json(
"orders",
[{"delete": {"id": 2, "customer": "bob", "amount": 5}}],
update_format="insert_delete",
)
rows = list(pipeline.query("SELECT * FROM customer_totals"))
pipeline.stop(force=True)
pipeline.clear_storage()
assert rows == [{"customer": "alice", "total": 25, "num_orders": 1}]
Deleting a row of a table without a primary key when the row is not in the table produces unpredictable results, and the pipeline may crash.
Check that a change does not occur
While Feldera has a reliable mechanism to track the completion of steps, called a completion token, unfortunately the HTTP connector (which is the basis for the listener) does not expose this information. We hope to improve this API: see issue 7360.
Some input changes produce no output changes (e.g., updating a record will
not change a row COUNT). To check that an input produces no output
changes, a test has to use a different approach:
- For materialized views, use an ad-hoc query executed after
input_json()returns. The snapshot will include all effects of the input. - For a non-materialized view, push a subsequent input which produces an expected output change. Since inputs are processed in order, receiving the change produced by the second input guarantees that the first one has been processed.
def check_identical_row_causes_no_change(pipeline: Pipeline) -> None:
# Connect the listener before the first step.
pipeline.start_paused()
totals = pipeline.listen("customer_totals")
pipeline.resume()
pipeline.input_json("orders", [{"id": 1, "customer": "alice", "amount": 10}])
# The input under test: replace order 1 with an identical row.
pipeline.input_json("orders", [{"id": 1, "customer": "alice", "amount": 10}])
# The second input: its effect is known.
pipeline.input_json("orders", [{"id": 2, "customer": "bob", "amount": 5}])
changes = read_changes(totals, 2)
pipeline.stop(force=True)
pipeline.clear_storage()
assert changes == [
{"customer": "alice", "total": 10, "num_orders": 1, "insert_delete": 1},
{"customer": "bob", "total": 5, "num_orders": 1, "insert_delete": 1},
]
read_changes(totals, 2) returns when the totals listener has
received at least two changes.
Tests with transactions
A transaction makes the pipeline process a group of inputs as one unit. The views change once, at commit, without exposing intermediate results.
def check_transaction_changes_view_once(pipeline: Pipeline) -> None:
# Connect the listener before the first step.
pipeline.start_paused()
totals = pipeline.listen("customer_totals")
pipeline.resume()
pipeline.start_transaction()
# The pipeline processes the input only at commit, so do not wait for it.
pipeline.input_json("orders", [{"id": 1, "customer": "alice", "amount": 10}], wait=False)
pipeline.input_json("orders", [{"id": 2, "customer": "alice", "amount": 7}], wait=False)
before_commit = list(pipeline.query("SELECT * FROM customer_totals"))
pipeline.commit_transaction()
changes = read_changes(totals, 1)
pipeline.stop(force=True)
pipeline.clear_storage()
assert before_commit == []
assert changes == [
{"customer": "alice", "total": 17, "num_orders": 2, "insert_delete": 1}
]
Compare this result with the test in Check the contents of a view as a change stream: the same two orders without a transaction cause three changes.
Before the commit, query() shows the contents of the views from
before the start of the transaction.
Testing programs that use the real-time clock NOW()
The NOW() function returns the current time of the system clock, in
UTC unless the clock_timezone_offset runtime setting gives another time
zone. This makes programs using NOW() non-deterministic: if the same
program is run twice over the same input, the results can be different.
As described above, Feldera pipelines are fully deterministic, and the
value of the clock is supplied to a pipeline using a compiler-synthesized
table which has a single column of type TIMESTAMP. By default a built-in
clock connector updates the sole value in this table periodically.
The following settings allow the test to set the value
returned by NOW():
| Setting | Effect |
|---|---|
dev_tweaks.now_offset | Initial value of the NOW() timestamp. |
dev_tweaks.now_http_driven | NOW() advances only when the test calls pipeline.advance_clock(delta_ms). |
clock_resolution_usecs | The step of the clock, in microseconds. pipeline.advance_clock() with no argument moves NOW() forward by this value. |
CREATE TABLE events (
id BIGINT NOT NULL,
t TIMESTAMP NOT NULL
);
-- Events of the last hour.
CREATE MATERIALIZED VIEW recent AS
SELECT id FROM events WHERE t >= NOW() - INTERVAL 1 HOUR;
-- The clock value of the latest step, in milliseconds since the epoch.
CREATE MATERIALIZED VIEW clock AS SELECT CAST(NOW() AS BIGINT) AS now_ms;
The settings are part of the runtime configuration, so the program is compiled with them:
from feldera.runtime_config import RuntimeConfig
CLOCK_CONFIG = RuntimeConfig(
clock_resolution_usecs=1_000_000,
dev_tweaks={"now_offset": "2030-01-01T00:00:00Z", "now_http_driven": True},
)
pipeline = PipelineBuilder(
client, name="test-recent", sql=SQL, runtime_config=CLOCK_CONFIG
).create_or_replace()
MINUTE_MS = 60_000
def clock_ms(pipeline: Pipeline) -> int | None:
rows = list(pipeline.query("SELECT now_ms FROM clock"))
return rows[0]["now_ms"] if rows else None
def advance_clock(pipeline: Pipeline, delta_ms: int) -> None:
"""Move NOW() forward and wait until all views use the new value."""
target_ms = pipeline.advance_clock(delta_ms)["now_ms"]
wait_for_condition(
f"NOW() is {target_ms} ms", lambda: clock_ms(pipeline) == target_ms
)
def check_event_leaves_window(pipeline: Pipeline) -> None:
pipeline.start()
pipeline.input_json("events", [{"id": 1, "t": "2030-01-01 00:00:00"}])
assert list(pipeline.query("SELECT id FROM recent")) == [{"id": 1}]
advance_clock(pipeline, 30 * MINUTE_MS)
assert list(pipeline.query("SELECT id FROM recent")) == [{"id": 1}]
advance_clock(pipeline, 31 * MINUTE_MS)
assert list(pipeline.query("SELECT id FROM recent")) == []
pipeline.stop(force=True)
pipeline.clear_storage()
advance_clock() returns before the pipeline completes a step with
the new value of NOW(). The clock view shows the new value only
after that step completes.
Tests with fixed inputs
The datagen connector
The datagen connector generates the same rows each time the pipeline starts, so it can be used to test programs with fixed input data. The following program generates 100 orders:
CREATE TABLE orders (
id BIGINT NOT NULL PRIMARY KEY,
customer VARCHAR NOT NULL,
amount INT NOT NULL
) WITH (
'connectors' = '[{
"transport": {
"name": "datagen",
"config": { "plan": [{ "limit": 100 }] }
}
}]'
);
CREATE MATERIALIZED VIEW customer_totals AS
SELECT customer, SUM(amount) AS total, COUNT(*) AS num_orders
FROM orders
GROUP BY customer;
The generator gives each column the values 0, 1, ..., 99, so each order has a different customer.
def check_generated_orders(pipeline: Pipeline) -> None:
pipeline.start()
pipeline.wait_for_completion()
rows = list(
pipeline.query(
"SELECT COUNT(*) AS customers, SUM(total) AS total FROM customer_totals"
)
)
pipeline.stop(force=True)
pipeline.clear_storage()
assert rows == [{"customers": 100, "total": 4950}]
After it generates the limit rows of its plan, the connector signals
the end of its input. Without a limit setting, the connector generates rows
forever, and wait_for_completion() never returns.
Reading S3 files
A test can also read a fixed input from a file in an object store. The following program uses the S3 connector to read a JSON file with three vendors from a public bucket:
CREATE TABLE vendor (
id BIGINT NOT NULL PRIMARY KEY,
name VARCHAR,
address VARCHAR
) WITH (
'connectors' = '[{
"transport": {
"name": "s3_input",
"config": {
"bucket_name": "feldera-basics-tutorial",
"key": "vendor.json",
"region": "us-west-1",
"no_sign_request": true
}
},
"format": { "name": "json" }
}]'
);
CREATE MATERIALIZED VIEW vendor_names AS SELECT id, name FROM vendor;
def check_vendors_from_s3(pipeline: Pipeline) -> None:
pipeline.start()
pipeline.wait_for_completion()
rows = list(pipeline.query("SELECT id, name FROM vendor_names ORDER BY id"))
pipeline.stop(force=True)
pipeline.clear_storage()
assert rows == [
{"id": 1, "name": "Gravitech Dynamics"},
{"id": 2, "name": "HyperDrive Innovations"},
{"id": 3, "name": "DarkMatter Devices"},
]
As with datagen, wait_for_completion() returns after the connector has
read the whole file.
Negative tests
Malformed programs
PipelineBuilder(...).create_or_replace() raises a RuntimeError when the program does
not compile. The message contains the errors of the compiler.
def check_unknown_column_is_rejected(client: FelderaClient) -> None:
sql = "CREATE TABLE t (x INT);\nCREATE VIEW v AS SELECT y FROM t;"
try:
PipelineBuilder(client, name="test-bad-sql", sql=sql).create_or_replace()
except RuntimeError as error:
assert "failed to compile" in str(error)
else:
raise AssertionError("the program compiled")
finally:
client.delete_pipeline("test-bad-sql")
To check only whether a program compiles, client.validate_program(sql)
is faster: it runs the SQL compiler, but it does not create a
pipeline. It does not raise an exception
for an error in the program. It returns {"Success": ...}, or
{"SqlError": ...} with one message for each error:
def check_unknown_column_is_rejected_quickly(client: FelderaClient) -> None:
sql = "CREATE TABLE t (x INT);\nCREATE VIEW v AS SELECT y FROM t;"
result = client.validate_program(sql)
messages = result["SqlError"]["info"]["messages"]
assert "Column 'y' not found in any table" in messages[0]["message"]
Runtime errors
Some errors happen only at runtime, for example an
arithmetic overflow. When a pipeline crashes, it
stops, and pipeline.deployment_error() describes the error.
CREATE TABLE numbers (a TINYINT NOT NULL, b TINYINT NOT NULL);
CREATE MATERIALIZED VIEW sums AS SELECT a + b AS total FROM numbers;
from feldera.enums import PipelineStatus
from feldera.rest.errors import FelderaAPIError
def check_overflow_stops_pipeline(pipeline: Pipeline) -> None:
pipeline.start()
try:
# 120 + 100 does not fit in a TINYINT.
pipeline.input_json("numbers", [{"a": 120, "b": 100}])
except FelderaAPIError:
pass
wait_for_condition(
"the pipeline stopped",
lambda: pipeline.status() == PipelineStatus.STOPPED,
timeout_s=60,
)
error = pipeline.deployment_error()
pipeline.stop(force=True)
pipeline.clear_storage()
assert error["error_code"] == "RuntimeError.WorkerPanic"
assert "'120 + 100' causes overflow for type TINYINT" in error["message"]
Unfortunately one cannot combine multiple checks for expected runtime in a single test: the pipeline will stop at the first error.
Common problems
| Symptom | Cause | Fix |
|---|---|---|
pipeline.query() fails for a view. | The view is not materialized. | Declare it CREATE MATERIALIZED VIEW, or read its change stream instead of using query. |
| The listener misses changes, or receives changes from before it connected. | The test connected the listener to a running pipeline. | Connect the listeners while the pipeline is paused: pipeline.start_paused(), pipeline.listen(), pipeline.resume(). |
| The listener has fewer changes than expected. | The test read the listener before all changes arrived. | Wait with read_changes(). |
| A test fails at random. | The test sleeps for a fixed time. | Synchronize as shown in Synchronization. |
| Tests fail at random when several tests run concurrently. | Pipeline names are unique within one account of a Feldera API. create_or_replace() stops a running pipeline with the same name and replaces it, so one run destroys the pipeline of another. | Give the pipelines of each run different names, for example by adding a suffix that is unique to the run. |
| A comparison fails, but the rows or changes are correct. | Their order is not defined. | Use ORDER BY in the query, or sorted_dicts(). |
A view with LIMIT and an ad-hoc query with the same LIMIT may return different rows. | When rows are tied on the ORDER BY columns, or there is no ORDER BY, LIMIT can keep any of the tied rows. The pipeline and ad-hoc queries use different SQL engines, which choose different rows. | Add columns to ORDER BY, for example the primary key. |
pipeline.input_json() does not return during a transaction. | The pipeline processes the input only at commit. | Invoke with wait=False. |
| A floating-point value does not compare equal. | pipeline.query() returns floating-point values as Python Decimal values, and rounding differs between engines. | Round the value, or CAST it to VARCHAR in the SQL view. |
| A value has an unexpected type. | Listeners return pandas values, for example pandas.Timestamp for TIMESTAMP. pipeline.query() returns JSON values, for example strings for TIMESTAMP. | Compare with values of the same type. |
| Column names do not match. | Feldera converts names that are not quoted to lowercase. | Use lowercase names in the expected rows. |
| A view shows wrong results after a delete. | The test deleted a row that a table without a primary key does not contain. | Delete only rows that the test inserted, or add a primary key. |
| A check passes although it should fail. | Python runs with -O, which removes assert statements. | Run the tests without -O. |
Avoiding hanging programs
The examples above wait without a time limit: each blocking call returns only when the pipeline has done its work. If the pipeline is stuck or does not even start, the test waits forever. A test that runs unattended, for example a unit test in a continuous integration job, must fail instead of hanging. Most blocking SDK calls accept a time limit:
| Call | Time limit | Error when the time is up |
|---|---|---|
pipeline.input_json() | wait_timeout_s=... | FelderaTimeoutError |
pipeline.start(), pipeline.start_paused(), pipeline.resume() | timeout_s=... | TimeoutError |
pipeline.commit_transaction() | timeout_s=... | TimeoutError |
pipeline.wait_for_completion() | timeout_s=... | TimeoutError |
pipeline.stop(), pipeline.clear_storage() | timeout_s=... | FelderaTimeoutError |
| Each HTTP request of a client | FelderaClient(url, timeout=...) | FelderaTimeoutError |
pipeline.input_pandas(), pipeline.execute(..., wait=True), compilation in PipelineBuilder(...).create_or_replace() | None |
The polling helper from Helper functions take a
limit argument: pass timeout_s to wait_for_condition() or to
read_changes(), for example read_changes(totals, 3, timeout_s=60).
They raise TimeoutError when the time is up.
from feldera.rest.errors import FelderaTimeoutError
TIMEOUT_S = 60
def check_totals_with_time_limits(pipeline: Pipeline) -> None:
pipeline.start(timeout_s=TIMEOUT_S)
try:
pipeline.input_json(
"orders",
[{"id": 1, "customer": "alice", "amount": 10}],
wait_timeout_s=TIMEOUT_S,
)
rows = list(pipeline.query("SELECT * FROM customer_totals"))
except (TimeoutError, FelderaTimeoutError):
print("status:", pipeline.status(), "errors:", pipeline.errors())
raise
finally:
pipeline.stop(force=True, timeout_s=TIMEOUT_S)
pipeline.clear_storage(timeout_s=TIMEOUT_S)
assert rows == [{"customer": "alice", "total": 10, "num_orders": 1}]
These methods can provide additional diagnostics:
| Operation | Shows |
|---|---|
pipeline.status() | The state of the pipeline, for example RUNNING or PAUSED |
pipeline.errors() | The compilation errors, and the error that stopped the pipeline, if any |
pipeline.stats() | The statistics, for example the number of input records received and processed |
pipeline.logs() | The log of the pipeline, one line at a time |