Skip to content
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
38 changes: 36 additions & 2 deletions docs.feldera.com/docs/connectors/sources/delta.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ exactly once fault tolerance.
| `mode`* | enum | | Table read mode. Four options are available: <ul> <li>`snapshot` - read a snapshot of the table and stop.</li> <li>`follow` - follow the changelog of the table, only ingesting changes (new and deleted rows)</li> <li>`snapshot_and_follow` - Read a snapshot of the table before switching to the `follow` mode. This mode implements the backfill pattern where we load historical data for the table before ingesting the stream of real-time updates.</li><li>`cdc` - Change-Data-Capture (CDC) mode. The connector treats the table as an append-only log where every row represents an insert or delete action. The order of actions is determined by the `cdc_order_by` property, and the type of each action is determined by the `cdc_delete_filter` property. Removed rows are ignored, so users can clean up old log entries without affecting the contents of the ingested stream. In this mode, the connector does not read the initial snapshot of the table and follows the transaction log starting from the version of the table specified by the `version` or `datetime` property.</li> </ul>|
| `transaction_mode` | enum | `none` | Determines how the connector breaks up its input into transactions. Supported values are `none`, `snapshot`, `catchup`, and `always`. See [below](#transactions) for details. |
| `timestamp_column` | string | | Table column that serves as an event timestamp. When this option is specified, and `mode` is one of `snapshot` or `snapshot_and_follow`, table rows are ingested in the timestamp order, respecting the [`LATENESS`](/sql/streaming#lateness-expressions) property of the column: each ingested row has a timestamp no more than `LATENESS` time units earlier than the most recent timestamp of any previously ingested row. See details [below](#ingesting-time-series-data-from-a-delta-lake). |
| `filter` | string | | <p>Optional row filter.</p> <p>When specified, only rows that satisfy the filter condition are read from the delta table. The condition must be a valid SQL Boolean expression that can be used in the `where` clause of the `select * from my_table where ...` query.</p> |
| <a name="filter">`filter`</a> | string | | <p>Optional row filter.</p> <p>When specified, only rows that satisfy the filter condition are read from the delta table. The condition must be a valid SQL Boolean expression that can be used in the `where` clause of the `select * from my_table where ...` query.</p> |
| `snapshot_filter` | string | | <p>Optional snapshot filter.</p><p>This option is only valid when `mode` is set to `snapshot` or `snapshot_and_follow`. When specified, only rows that satisfy the filter condition are included in the snapshot.</p> <p>The condition must be a valid SQL Boolean expression that can be used in the `where` clause of the `select * from snapshot where ...` query.</p><p>Unlike the `filter` option, which applies to all records retrieved from the table, this filter only applies to rows in the initial snapshot of the table. For instance, it can be used to specify the range of event times to include in the snapshot, e.g.: `ts BETWEEN TIMESTAMP '2005-01-01 00:00:00' AND TIMESTAMP '2010-12-31 23:59:59'`. This option can be used together with the `filter` option. During the initial snapshot, only rows that satisfy both `filter` and `snapshot_filter` are retrieved from the Delta table. When subsequently following changes in the the transaction log (`mode = snapshot_and_follow`), all rows that meet the `filter` condition are ingested, regardless of `snapshot_filter`. </p> |
| `version`, `start_version` | integer| | <p>Optional table version. When this option is set, the connector finds and opens the specified version of the table. In `snapshot` and `snapshot_and_follow` modes, it retrieves the snapshot of this version of the table. In `follow`, `snapshot_and_follow`, and `cdc` modes, it follows transaction log records **after** this version.</p><p>Note: at most one of `version` and `datetime` options can be specified. When neither of the two options is specified, the latest committed version of the table is used.</p> |
| `datetime` | string | | <p>Optional timestamp for the snapshot in the ISO-8601/RFC-3339 format, e.g., "2024-12-09T16:09:53+00:00". When this option is set, the connector finds and opens the version of the table as of the specified point in time (based on the server time recorded in the transaction log, not the event time encoded in the data). In `snapshot` and `snapshot_and_follow` modes, it retrieves the snapshot of this version of the table. In `follow`, `snapshot_and_follow`, and `cdc` modes, it follows transaction log records **after** this version.</p><p> Note: at most one of `version` and `datetime` options can be specified. When neither of the two options is specified, the latest committed version of the table is used.</p>|
Expand Down Expand Up @@ -359,6 +359,40 @@ it is marked **UNHEALTHY** while retrying failed operations.
If the pipeline is stopped and restarted during a retry, the connector resumes from the last successfully
ingested table version. This guarantees that no data loss occurs due to object store read errors.

## Optimizing multihost performance for unbalanced large data

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

much of this is not specific to Delta, maybe there should be a separate page about this.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is definitely true as the amount of advice increases; I wasn't sure whether it was yet for n = 2.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This advice doesn't work for tables with lateness, since different connectors can go out of sync.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I fixed this by adding a warning paragraph:

> ⚠️ [Tables with LATENESS] cannot correctly be divided into multiple
> connectors this way, because the different connectors can read data
> "out of sync" from another.

[Tables with LATENESS]: /sql/streaming/#lateness-expressions

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

perfect


In a multihost pipeline, Feldera currently assigns each input
connector to one host. All of the work related to reading data from
that input connector then takes place on that host. When a pipeline
has several input connectors that read comparable amounts of data,
this spreads the network and CPU load related to them across the
hosts.

On the other hand, if a pipeline that has one input connector that
reads most of the pipeline's data, only a single host does all the
work. In such a case, it makes sense to divide the connector into
multiple connectors, one per host. Feldera places input connectors on
hosts "round robin" in alphabetical order, which means that using the
default connector names or names with the same prefix and a sequential
suffix will ensure that they are spread as evenly as possible across
the hosts.

When a Delta Lake source is divided into
[partitions](https://docs.delta.io/best-practices/), its input
connector can be divided into multiple connectors by specifying a
different [`filter`](#filter) for each one. For example, if the
source is partitioned by an integer column `partition` that takes one
Comment thread
blp marked this conversation as resolved.
of the six values 0 through 5, and there are 3 hosts, one might use 3
copies of the input connector, one with `"filter": "partition = 0 OR
partition = 1"`, one with `"filter": "partition = 2 OR partition =
3"`, and one with `"filter": "partition = 4 OR partition = 5`.

> ⚠️ [Tables with LATENESS] cannot correctly be divided into multiple
> connectors this way, because the different connectors can read data
> "out of sync" from one another.

[Tables with LATENESS]: /sql/streaming/#lateness-expressions

## Additional examples

### Example: Setting `timestamp_column`
Expand Down Expand Up @@ -441,4 +475,4 @@ Read table snapshot via Unity catalog.
"unity_host": "https://dbc-XXX-XXX.cloud.databricks.com",
"aws_region": "us-west-1"
}
```
```
38 changes: 36 additions & 2 deletions docs.feldera.com/docs/connectors/sources/kafka.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ tolerance](/pipelines/fault-tolerance).
| `topic` (required) | string | | The Kafka topic to subscribe to. |
| `bootstrap.servers` (required) | string | | A comma separated list of Kafka brokers to connect to.
| `start_from` | variant | latest | The starting point for reading from the topic, one of `"earliest"`, `"latest"`, `{"offsets": [1, 2, 3, ...]}`, where `[1, 2, 3, ...]` are particular offsets within each partition in the topic, or `{"timestamp": <timestamp>}` where `<timestamp>` is a Kafka timestamp as an integer number of milliseconds since the epoch. |
| `partitions` | integer list | | <p> The list of Kafka partitions to read from. </p> <p> Only the specified partitions will be consumed. If this field is not set, the connector will consume from all available partitions. </p><p> If `start_from` is set to `offsets` and this field is provided, the number of partitions must exactly match the number of offsets, and the order of partitions must correspond to the order of offsets. </p><p> If offsets are provided for all partitions, this field can be omitted. </p> |
| <a name="partitions">`partitions`</a> | integer list | | <p> The list of Kafka partitions to read from. </p> <p> Only the specified partitions will be consumed. If this field is not set, the connector will consume from all available partitions. </p><p> If `start_from` is set to `offsets` and this field is provided, the number of partitions must exactly match the number of offsets, and the order of partitions must correspond to the order of offsets. </p><p> If offsets are provided for all partitions, this field can be omitted. </p> |
| `log_level` | string | | The log level for the Kafka client. |
| `group_join_timeout_secs` | seconds | 10 | Maximum timeout (in seconds) for the endpoint to join the Kafka consumer group during initialization. |
| `poller_threads` | positive integer | 3 | Number of threads used to poll Kafka messages. Setting it to multiple threads can improve performance with small messages. Default is 3. Each partition being read is assigned to one of the threads, so the connector automatically caps the thread count at the partition count. |
Expand All @@ -23,7 +23,7 @@ tolerance](/pipelines/fault-tolerance).
| `include_partition` | boolean | false | Whether to include Kafka partition in connector metadata (see [Accessing Kafka metadata](#metadata)). |
| `include_offset` | boolean | false | Whether to include Kafka offset name in connector metadata (see [Accessing Kafka metadata](#metadata)). |
| `include_timestamp` | boolean | false | Whether to include Kafka timestamp in connector metadata (see [Accessing Kafka metadata](#metadata)). |
| `synchronize_partitions` | boolean | false | Whether to read records in order of Kafka timestamp across partitions (see [Synchronizing partitions](#synchronizing-partitions)) |
| <a name="synchronize_partitions">`synchronize_partitions`</a> | boolean | false | Whether to read records in order of Kafka timestamp across partitions (see [Synchronizing partitions](#synchronizing-partitions)) |
| `header_filter` | filter | | Drop messages whose Kafka headers do not match a boolean expression of regular expressions (see [Filtering messages by header](#header-filter)). |

The connector passes additional options directly to [**librdkafka**](https://github.com/confluentinc/librdkafka/blob/master/CONFIGURATION.md). Some of the relevant options:
Expand Down Expand Up @@ -508,6 +508,40 @@ Pitfalls of this solution include:
be added to the partition whose events have been completely
processed.

## Optimizing multihost performance for unbalanced large data

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Kafka already doesn't guarantee lateness across multiple partitions, but we have the synchronize_partitions feature, which is not applicable if we split ingest across multiple connectors.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I fixed this by adding a warning paragraph:

> ⚠️ The [synchronize_partitions] feature does not work across
> connectors, only within a single connector.  Therefore, [tables with
> LATENESS] cannot correctly be divided into multiple connectors this
> way, because the different connectors can read data "out of sync"
> from another.

[tables with LATENESS]: /sql/streaming/#lateness-expressions
[synchronize_partitions]: #synchronize_partitions

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sounds good


In a multihost pipeline, Feldera currently assigns each input
connector to one host. All of the work related to reading data from
that input connector then takes place on that host. When a pipeline
has several input connectors that read comparable amounts of data,
this spreads the network and CPU load related to them across the
hosts.

On the other hand, if a pipeline that has one input connector that
reads most of the pipeline's data, only a single host does all the
work. In such a case, it makes sense to divide the connector into
multiple connectors, one per host. Feldera places input connectors on
hosts "round robin" in alphabetical order, which means that using the
default connector names or names with the same prefix and a sequential
suffix will ensure that they are spread as evenly as possible across
the hosts.

A Kafka input connector can be divided into multiple connectors by
specifying different [`partitions`](#partitions) for each one. For
Comment thread
blp marked this conversation as resolved.
example, if the Kafka topic has 6 partitions, and there are 3 hosts,
one might use 3 copies of the input connector, one with `"partitions":
[0, 1]`, one with `"partitions": [2, 3]`, and one with `"partitions":
[4, 5]`.

> ⚠️ The [synchronize_partitions] feature does not work across
> connectors, only within a single connector. Therefore, [tables with
> LATENESS] cannot correctly be divided into multiple connectors this
> way, because the different connectors can read data "out of sync"
> from one another.

[tables with LATENESS]: /sql/streaming/#lateness-expressions
[synchronize_partitions]: #synchronize_partitions

## Additional resources

For more information, see:
Expand Down
Loading