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
9 changes: 8 additions & 1 deletion doc/user/content/concepts/snapshotting.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,10 @@ menu:

{{% include-headless "/headless/ingestion/snapshotting-duration" %}}

### Parallelism

{{% include-headless "/headless/ingestion/snapshotting-parallelism" %}}

## Queries during snapshotting

{{% include-headless "/headless/ingestion/snapshotting-queries" %}}
Expand All @@ -27,7 +31,9 @@ menu:
Snapshotting has the following upstream impacts:

- **Read load.** Snapshotting puts read, CPU, and network load on the upstream
system, proportional to the data volume.
system. The total load is proportional to the volume of data being
snapshotted, while the source cluster's [parallelism](#parallelism) affects
the peak load: more workers compress the reads into a shorter window.

- **Change-log retention for CDC database sources.** When ingesting data from
CDC database sources (PostgreSQL, MySQL, SQL Server), the upstream system must
Expand All @@ -41,3 +47,4 @@ Snapshotting has the following upstream impacts:

- [Ingest data](/ingest-data/)
- [Sources](/concepts/sources/)
- [Troubleshooting data ingestion](/ingest-data/troubleshooting/)
34 changes: 34 additions & 0 deletions doc/user/content/headless/ingestion/snapshotting-parallelism.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
---
headless: true
---

Materialize can parallelize snapshotting across the workers of the cluster
hosting the source.

- **PostgreSQL sources** are parallelized by table, i.e., different tables
are read concurrently by different workers. On PostgreSQL 14 and later,
Materialize additionally attempts to partition each table's read across
workers. Tables that cannot be partitioned fall back to a single worker.

- **MySQL sources** are parallelized by table, i.e., different tables are
read concurrently by different workers. For tables that meet certain
requirements, Materialize can additionally partition the table's read
across workers {{< private-preview-inline />}}. See [MySQL snapshot
parallelism](/ingest-data/mysql/snapshot-parallelism/).

- **Kafka sources** are parallelized by topic partition, with partitions
distributed across workers, so parallelism is bounded by the topic's
partition count.

- **SQL Server sources** are not parallelized: a single worker reads all
tables.

The degree of snapshot parallelism depends on the number of workers. A
cluster's [size](/sql/create-cluster/#available-sizes) determines its number
of workers, so a larger cluster can shorten the snapshot, to the extent the
work parallelizes and the upstream database keeps up. The volume read from
the upstream database is unchanged, it is compressed into a shorter window
of more concurrent queries and connections. To determine whether
snapshotting is overloading the upstream database, and for ways to mitigate
the load, see [Is the upstream database
overloaded?](/ingest-data/troubleshooting/#is-the-upstream-database-overloaded)
4 changes: 4 additions & 0 deletions doc/user/content/ingest-data/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,10 @@ we recommend:
the steady-state resource needs of your upsert source(s). See [Best practices:
Upsert sources](#upsert-sources).

### Parallelism

{{% include-headless "/headless/ingestion/snapshotting-parallelism" %}}

### Monitoring progress

While snapshotting is taking place, you can monitor the progress of the
Expand Down
5 changes: 5 additions & 0 deletions doc/user/content/ingest-data/mysql/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,11 @@ gives you the following benefits:
read-replica to build views on top of your MySQL data that are efficiently
maintained and always up-to-date.

When a source is created, Materialize parallelizes the initial snapshot
across the cluster's workers and can split the read of large tables that meet
certain requirements {{< private-preview-inline />}}. See [Snapshot
parallelism](/ingest-data/mysql/snapshot-parallelism/).

## Supported versions and services

{{< note >}}
Expand Down
108 changes: 108 additions & 0 deletions doc/user/content/ingest-data/mysql/snapshot-parallelism.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
---
title: "Snapshot parallelism"
description: "How Materialize splits the snapshot of a single MySQL table across the workers of a cluster."
menu:
main:
parent: "mysql"
Comment thread
kay-kim marked this conversation as resolved.
name: "Snapshot parallelism"
identifier: "mysql-snapshot-parallelism"
weight: 70
---

{{< private-preview />}}

When you create a [MySQL source](/sql/create-source/mysql-v2/), Materialize
Comment thread
kay-kim marked this conversation as resolved.
performs an initial, snapshot-based sync of the selected tables before it
starts ingesting change events from the binlog. For large tables, this
snapshot dominates the time until the source becomes healthy.

How snapshot work is spread across the workers of a cluster, and what that
means for the upstream database, is covered in
[Snapshotting](/concepts/snapshotting/#parallelism). Materialize can split
the read of a **single table** across all the workers of the cluster, so
that even a source dominated by one very large table benefits from a larger
cluster. This page covers what is specific to MySQL: which tables are
eligible for splitting, and how their reads are partitioned.

## Which tables are split

Materialize splits the snapshot of an individual table across workers when
all of the following conditions are met:

- The table has a **single-column primary key**. Composite primary keys are
not supported.
- The primary key column is of type **`CHAR` or `VARCHAR`**, with a declared
length of **at most 768 characters**. Other types, including numeric keys,
are not supported.
- The primary key column uses the **`utf8mb4` character set** with the
**`utf8mb4_bin` collation**.
- The table is **large enough to be worth splitting**. Small tables are read
by a single worker, where splitting would add overhead without benefit.

How evenly the split lands also depends on the distribution of the key
values. See [How a table is partitioned](#how-a-table-is-partitioned).

If a table does not meet these requirements, or if the [boundary
sampling](#how-a-table-is-partitioned) fails, its snapshot is not split: a
single worker reads the table in full. Different tables are still read
concurrently by different workers.

## How a table is partitioned

Materialize partitions an [eligible](#which-tables-are-split) table using the
leading characters of its primary key values. Before reading the table,
Materialize probes the primary key index to discover key prefixes and uses
the MySQL optimizer's row estimates to gauge how many rows fall under each
prefix. It extends the prefixes as needed to find boundaries that divide the
table into roughly even ranges. The probes are inexpensive point lookups,
capped in proportion to the table's estimated size, so the sampling phase
stays negligible next to the snapshot itself.

Each worker then reads only its assigned range, within the same consistent
snapshot of the upstream database, so the result is identical to a
single-worker snapshot, only faster.

Because partitioning is based on key prefixes and optimizer estimates, how
evenly the work divides depends on the shape of your keys:

- **Evenly distributed keys partition well.** Keys whose leading characters
spread rows uniformly, such as UUIDs, hashes, or other randomized
identifiers, produce well-balanced ranges.

- **Skewed keys partition less evenly.** If a large share of the table's rows
sort under a few common prefixes, some ranges end up with more rows than
others, and the workers assigned to them finish later.

- **The probe budget can run out.** If finding even boundaries would require
examining very many distinct prefixes, Materialize stops probing and uses
the coarser boundaries found so far, which can also leave ranges uneven.

Uneven partitioning is never incorrect. It only reduces the speedup, since
the snapshot finishes when the busiest worker finishes.

## MySQL-specific upstream considerations

- **Connection count.** While the snapshot is being set up, Materialize
briefly holds up to two connections per worker, plus one. Once reading is
underway, this settles to one connection per worker reading a range, plus
one coordination connection. After the snapshot completes, the source drops
back to a single replication connection. If your MySQL server or connection
pooler enforces a low
[`max_connections`](https://dev.mysql.com/doc/refman/8.0/en/server-system-variables.html#sysvar_max_connections)
limit, account for this burst when sizing it.

- **Statistics freshness.** Range boundaries are placed using the MySQL
optimizer's row estimates. Stale statistics don't affect correctness, but
can skew how evenly work divides across workers. Running
[`ANALYZE TABLE`](https://dev.mysql.com/doc/refman/8.0/en/analyze-table.html)
on very large tables before creating the source can improve balance.
Comment on lines +84 to +98

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I'd keep these here for now: both bullets only apply when parallel snapshotting is active, which is private preview and flag-off, so they'd be noise in the general MySQL considerations. Worth revisiting when the feature is on by default.


For general guidance on read load, IOPS, and other upstream impact, which is
not specific to MySQL, see [Is the upstream database
overloaded?](/ingest-data/troubleshooting/#is-the-upstream-database-overloaded)

## Observability

To observe the progress of an ongoing snapshot, see [Monitoring the
snapshotting
progress](/ingest-data/monitoring-data-ingestion/#monitoring-the-snapshotting-progress).
5 changes: 5 additions & 0 deletions doc/user/content/ingest-data/postgres/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,11 @@ Materialize gives you the following benefits:
Materialize as a read-replica to build views on top of your PostgreSQL data
that are efficiently maintained and always up-to-date.

When a source is created, Materialize parallelizes the initial snapshot
across the cluster's workers and, on PostgreSQL 14 and later, splits each
table's read across workers. See [Snapshot
parallelism](/concepts/snapshotting/#parallelism).

## Supported versions and services

The PostgreSQL source requires **PostgreSQL 11+** and is compatible with most
Expand Down
35 changes: 35 additions & 0 deletions doc/user/content/ingest-data/troubleshooting.md
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,41 @@ also be necessary to support increased memory usage during the process. For more
information, see [Use a larger cluster for upsert source
snapshotting](/ingest-data/#use-a-larger-cluster-for-upsert-source-snapshotting).

## Is the upstream database overloaded?

Snapshotting can put significant load on the upstream database (see [Impact
on upstream system](/concepts/snapshotting/#impact-on-upstream-system)).

Check the upstream database when a snapshot progresses more slowly than
expected, when applications sharing the database slow down while
it runs, or when the source reports upstream connection errors or timeouts.
The relevant metrics are in your cloud provider's monitoring console, or in
OS tools like `iostat` and the database's activity views for self-hosted
databases. Look for:

- **Read IOPS or throughput** flat at a provisioned cap.
- **CPU** pinned at the instance's limit for the duration of the snapshot.
- **Network throughput** at the instance type's cap.
- **Connections** near the database's limit. For PostgreSQL and MySQL
sources, snapshotting opens connections in proportion to the source
cluster's workers.

Also watch disk usage on the upstream database during a long-running
snapshot: CDC database sources must retain their change log until Materialize
consumes it (see [Impact on upstream
system](/concepts/snapshotting/#impact-on-upstream-system)).

If the database is overloaded, you can upsize the source database or cancel
the snapshot by dropping the source, and retry:

- on a smaller source cluster to spread the load over a longer window.
- with more IOPS, throughput, or instance capacity provisioned for the
database.
- during off-peak hours when the database is less busy, as recommended in the
[ingestion best practices](/ingest-data/#scheduling).
- with a smaller [volume of data to
sync](/ingest-data/#limit-the-volume-of-data).

## Adding a new subsource to an existing source blocks replication. Should I just create a new source instead?

It depends. Materialize provides transactional guarantees for subsource of the
Expand Down
Loading