Consume from Snowflake on AWS S3 Tables

View as Markdown

Once Materialize writes to an Iceberg table on Amazon S3 Tables, Snowflake can query that table directly through a catalog integration. Materialize and Snowflake never communicate with each other: both are clients of the Iceberg catalog, which holds the only shared state. Materialize commits snapshots to the catalog, and Snowflake polls the catalog for new snapshots.

The steps in this guide are specific to Iceberg tables hosted on Amazon S3 Tables. The catalog integration type, authentication method, and IAM permissions differ for Iceberg tables hosted elsewhere.

! Important: Create the sink with MODE APPEND. Snowflake cannot read Iceberg equality delete files, which MODE UPSERT writes whenever a row from an earlier snapshot is updated or deleted. For details, see Sink mode requirement for Snowflake.

Prerequisites

  • An Iceberg sink writing to Amazon S3 Tables, set up by following the AWS S3 Tables guide, created with MODE APPEND.

  • An AWS account with permissions to create and manage IAM policies and roles. This is a second IAM role, separate from the one Materialize assumes to write.

  • A Snowflake account, and a role with the global CREATE INTEGRATION privilege (an account-level privilege, typically held by ACCOUNTADMIN) and the CREATE DATABASE privilege.

Sink mode requirement for Snowflake

Snowflake’s support for Iceberg tables it does not manage itself excludes equality delete files: row-level deletes with equality delete files aren’t supported. This determines which sink mode you can use:

Sink mode Delete files written Readable by Snowflake
MODE APPEND None. Changes are data rows tagged with _mz_diff and _mz_timestamp. Yes
MODE UPSERT Equality deletes, for any update or delete against a row from an earlier snapshot. No

With MODE UPSERT, Snowflake fails in one of two ways depending on when the table was registered:

  • If the table was registered before the first equality delete was written, Snowflake stops refreshing the table and continues serving the last snapshot it could read. Queries succeed and return stale data without reporting an error. Once refresh stops, no further changes appear, including inserts that involve no delete files at all.

  • If the table was registered after equality deletes exist, Snowflake never reads any snapshot, and queries fail with 091968 (0A000): Equality deletes on Iceberg tables are not supported.

Materialize reports an upsert sink in this state as running with no error, because Materialize is writing valid Iceberg. The incompatibility is entirely on the read side. To detect it, monitor from Snowflake, as described in Monitor Snowflake’s refresh state.

Step 1. Create an IAM role for Snowflake

Create an IAM policy that grants read access to your S3 Tables catalog, replacing <S3 table bucket ARN> with the ARN of your S3 table bucket:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "s3tables:GetTableBucket",
                "s3tables:ListNamespaces",
                "s3tables:GetNamespace",
                "s3tables:ListTables",
                "s3tables:GetTable",
                "s3tables:GetTableData",
                "s3tables:GetTableMetadataLocation"
            ],
            "Resource": [
                "<S3 table bucket ARN>",
                "<S3 table bucket ARN>/table/*"
            ]
        }
    ]
}

Then create an IAM role that Snowflake can assume, and attach the policy to it. For the Trusted entity type, specify Custom trust policy. Both the principal and the external ID are placeholders at this stage: Snowflake generates the real values when you create the catalog integration in the next step.

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Principal": {
                "AWS": "arn:aws:iam::<your account ID>:root"
            },
            "Action": "sts:AssumeRole",
            "Condition": {
                "StringEquals": {
                    "sts:ExternalId": "PENDING"
                }
            }
        }
    ]
}

Note down the role ARN. You will use it in the next step.

Step 2. Create the catalog integration in Snowflake

In Snowflake, create a catalog integration for the S3 Tables Iceberg REST endpoint, replacing:

  • <region> with the AWS region of your S3 table bucket (for example, us-east-1),
  • <S3 table bucket ARN> with your S3 table bucket ARN, and
  • <Snowflake IAM role ARN> with the role ARN from step 1.
CREATE CATALOG INTEGRATION mz_s3tables_catalog
  CATALOG_SOURCE = ICEBERG_REST
  TABLE_FORMAT = ICEBERG
  REST_CONFIG = (
    CATALOG_URI = 'https://s3tables.<region>.amazonaws.com/iceberg'
    CATALOG_API_TYPE = AWS_S3TABLES
    ACCESS_DELEGATION_MODE = VENDED_CREDENTIALS
    CATALOG_NAME = '<S3 table bucket ARN>'
  )
  REST_AUTHENTICATION = (
    TYPE = SIGV4
    SIGV4_IAM_ROLE = '<Snowflake IAM role ARN>'
    SIGV4_SIGNING_REGION = '<region>'
  )
  ENABLED = TRUE;

ACCESS_DELEGATION_MODE = VENDED_CREDENTIALS means S3 Tables issues Snowflake scoped credentials for the underlying storage, so no external volume is required.

Next, retrieve the IAM principal Snowflake uses to assume your role:

DESCRIBE CATALOG INTEGRATION mz_s3tables_catalog;

Note down the values of the API_AWS_IAM_USER_ARN and API_AWS_EXTERNAL_ID properties. You will use them in the next step.

Step 3. Update the IAM role’s trust policy in AWS

In AWS, edit the trust policy of the IAM role you created in step 1, replacing the placeholder principal and external ID with the API_AWS_IAM_USER_ARN and API_AWS_EXTERNAL_ID values from the previous step:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Principal": {
                "AWS": "arn:aws:iam::781425929845:user/xxxxxxxx-x"
            },
            "Action": "sts:AssumeRole",
            "Condition": {
                "StringEquals": {
                    "sts:ExternalId": "XXXXXXXX_SFCRole=NN_xxxxxxxxxxxxxxxxxxxxxxxxxxxx="
                }
            }
        }
    ]
}

Step 4. Create a catalog-linked database in Snowflake

A catalog-linked database discovers the namespaces and tables in your Iceberg catalog and keeps them registered as they change, including tables Materialize creates later. Replace <namespace> with the namespace your sink writes to:

CREATE DATABASE mz_iceberg
  LINKED_CATALOG = (
    CATALOG = 'mz_s3tables_catalog',
    ALLOWED_NAMESPACES = ('<namespace>'),
    SYNC_INTERVAL_SECONDS = 30
  );

Discovery takes up to one SYNC_INTERVAL_SECONDS cycle. To confirm the tables have been registered:

SHOW ICEBERG TABLES IN DATABASE mz_iceberg;

Each table your sink writes to should appear with iceberg_table_type set to UNMANAGED.

Step 5. Reconstruct current state in Snowflake

In append mode, the Iceberg table is a changelog rather than a snapshot of current state, with two metadata columns added by the sink: _mz_diff and _mz_timestamp. Every change is written as a data row: _mz_diff is +1 for an insertion and -1 for a deletion, and an update appears as both rows, sharing one _mz_timestamp.

A query engine reading the table can reconstruct current state in one of two ways:

Approach How it works When to use it
Consolidate by diff Group by every column of the Iceberg table except _mz_diff and _mz_timestamp, and keep the groups whose _mz_diff values sum to a positive number. The sinked relation has no unique key, or you want a query that does not depend on one.
Latest version per key Rank the rows within each key by _mz_timestamp descending, breaking ties on _mz_diff descending, then keep the top-ranked row where _mz_diff is +1. The sinked relation has a unique key. Avoids grouping by every column, so it does not grow harder to write as the relation gets wider.

Both approaches return the same result. Two details matter for correctness:

  • An update writes -1 and +1 at the same _mz_timestamp, so ranking by timestamp alone is ambiguous. Breaking ties on _mz_diff descending is what selects the new version of the row rather than the old one.

  • The latest row for a deleted key carries _mz_diff = -1, so filtering for +1 after ranking is what removes deleted keys from the result.

Identifier quoting rules differ between query engines. Materialize creates Iceberg identifiers in lowercase unless they were quoted when created, so engines that resolve unquoted identifiers as uppercase require those identifiers to be quoted.

The examples assume that the sink writes a relation with columns id (the unique key), name, qty, price, and updated_at.

In Snowflake, quote the column and table names as lowercase. Snowflake resolves unquoted identifiers as uppercase, which will not match the identifiers Materialize created.

Group by every column except _mz_diff and _mz_timestamp, and keep the groups whose _mz_diff values sum to a positive number:

SELECT "id", "name", "qty"
  FROM mz_iceberg."<namespace>"."<table>"
 GROUP BY "id", "name", "qty", "price", "updated_at"
HAVING SUM("_mz_diff") > 0;

Every column of the table except _mz_diff and _mz_timestamp must appear in the GROUP BY clause, including the columns you do not select. Grouping on a subset merges rows that differ in the omitted columns and produces incorrect results.

Rank the rows within each key, then keep the top-ranked row for each key where _mz_diff is +1, replacing "id" with the unique key of your relation:

SELECT "id", "name", "qty"
  FROM (
    SELECT "id", "name", "qty", "_mz_diff",
           ROW_NUMBER() OVER (
             PARTITION BY "id"
             ORDER BY "_mz_timestamp" DESC, "_mz_diff" DESC) AS rn
      FROM mz_iceberg."<namespace>"."<table>"
  )
 WHERE rn = 1 AND "_mz_diff" = 1;

Ordering by "_mz_timestamp" DESC, "_mz_diff" DESC selects the new version of an updated row, and the "_mz_diff" = 1 filter removes keys whose most recent change was a deletion.

To define current state once and query it by name, wrap either query in a view. See Query cost of the changelog in Snowflake for how to bound the cost of doing so as the changelog grows.

To verify the pipeline end to end, compare the result against the same relation in Materialize. Compare values rather than row counts alone: a table that has stopped refreshing can hold a plausible number of rows whose values are stale.

Considerations

End-to-end latency from Materialize to Snowflake

Changes become visible in Snowflake in batches, not continuously. Two intervals compose:

  • The sink’s COMMIT INTERVAL determines when a new Iceberg snapshot exists at all. Every change written within one interval becomes visible at the same moment.
  • The catalog-linked database’s SYNC_INTERVAL_SECONDS (default 30) determines how soon after that Snowflake notices the new snapshot.

With COMMIT INTERVAL = '1m' and the default sync interval, a row’s visibility delay therefore depends on where in the commit window it was written, ranging from a few seconds to roughly the sum of both intervals.

Lowering COMMIT INTERVAL reduces latency at a cost:

The COMMIT INTERVAL setting controls how frequently Materialize commits snapshots to your Iceberg table, making the data available to downstream query engines. This setting involves tradeoffs:

Shorter intervals (e.g., < 1m) Longer intervals (e.g., 5m)
Lower latency - data visible sooner in downstream systems Higher latency - data takes longer to appear
More small files - can degrade query performance over time Fewer, larger files - better query performance
More frequent snapshot commits - higher catalog overhead Less catalog overhead
Lower throughput efficiency Higher throughput efficiency

Recommendations:

  • For production, use intervals of 1m or longer
  • For batch analytics, use longer intervals (5m to 15m)

Starting in v26.34, you can change the commit interval of an existing sink with ALTER SINK.

NOTE: Outside of development environments, commit intervals should be at least 1m. Short commit intervals increase catalog overhead and produce many small files. Small files will result in degraded query performance. It also increases load on the Iceberg metadata, which can result in a degraded catalog, and non-responsive queries.

Query cost of the changelog in Snowflake

The changelog grows with every change, so the query in step 5 reads more data over time. Defining it as a view keeps it always current but recomputes the full aggregation on each query. To bound query cost instead, materialize the result as a dynamic table, which adds its own refresh lag and compute cost. Dynamic tables track changes to externally managed Iceberg tables at file level, which suits an append-only source because appends add files without rewriting existing ones.

Monitor Snowflake’s refresh state

A query against a table that has stopped refreshing succeeds and returns the last snapshot Snowflake could read, so query results do not indicate whether refresh is healthy. Materialize’s sink status does not either. Check Snowflake’s refresh state instead:

SHOW ICEBERG TABLES IN DATABASE mz_iceberg;

In the auto_refresh_status column, executionState is RUNNING when refresh is healthy. Any other state includes an invalidExecutionStateReason explaining why refresh stopped, along with the metadata file and snapshot ID that failed. Snowflake also provides SYSTEM$AUTO_REFRESH_STATUS and ICEBERG_TABLE_SNAPSHOT_REFRESH_HISTORY for this purpose.

Automated refresh polls the catalog rather than relying on notifications, and is billed as Snowpipe usage.

AWS and Snowflake region placement

Your S3 table bucket must be in the same AWS region as your Materialize deployment. Placing your Snowflake account in that region as well avoids cross-region data transfer costs and added latency when Snowflake reads the table’s data files.

Troubleshooting

Snowflake queries return data older than expected

Refresh has most likely stopped. Check executionState as described in Monitor Snowflake’s refresh state. If invalidExecutionStateReason reports unsupported equality deletes, the sink is running in upsert mode; recreate it with MODE APPEND, as described in Sink mode requirement for Snowflake.

Forcing a refresh surfaces the same underlying error as a query error, which can be useful when diagnosing:

ALTER ICEBERG TABLE mz_iceberg."<namespace>"."<table>" REFRESH;

Snowflake reports that equality deletes are not supported

Snowflake cannot read the table because the sink writes equality delete files. Recreating the catalog integration or the catalog-linked database does not resolve this, because the delete files are in the table itself. Recreate the sink with MODE APPEND writing to a new Iceberg table.

Tables do not appear in the Snowflake catalog-linked database

Confirm that the namespace is listed in ALLOWED_NAMESPACES, that the sink has committed at least once, and that at least one SYNC_INTERVAL_SECONDS cycle has elapsed. If tables are still missing, verify the IAM role’s trust policy matches the API_AWS_IAM_USER_ARN and API_AWS_EXTERNAL_ID reported by DESCRIBE CATALOG INTEGRATION.

Back to top ↑