All skills
aws avatar

/ingesting-into-data-lake

@b33847d

Import data into the AWS data lake from S3 files, local uploads, JDBC databases (Oracle, SQL Server, PostgreSQL, MySQL, RDS, Aurora), Amazon Redshift, Snowflake, BigQuery, DynamoDB, or existing Glue catalog tables (migration). Default target is S3 Tables; standard Iceberg on a general purpose bucket is supported where S3 Tables is not adopted. Handles one-time loads, recurring pipelines, migrations. Triggers on: import data, load data, ingest, sync database, migrate table, move data to AWS, set up pipeline, ETL, pull from Snowflake, query BigQuery into S3, export DynamoDB, CTAS, convert to Iceberg. Do NOT use for setting up or troubleshooting Glue connections (use connecting-to-data-source), creating empty tables (use creating-data-lake-table), running queries (use querying-data-lake), finding tables by fuzzy name (use finding-data-lake-assets), catalog audit (use exploring-data-catalog), or SaaS platforms like Salesforce, ServiceNow, SAP, MongoDB, Kafka.

Use this Skill: https://skilld.dev/gh/aws/agent-toolkit-for-aws/ingesting-into-data-lake

This session only. Nothing lands on disk.

referencessnowflake-ingest.md

≈933 tokens on demand. Your agent reads this file only when SKILL.md points to it.

Snowflake Ingest

Move data from Snowflake into the data lake. Assumes a Glue SNOWFLAKE connection exists. If not, delegate to connecting-to-data-source.

Contents

Prerequisites

  • Glue connection of type SNOWFLAKE (not JDBC)
  • Source database, schema, table, and optional query
  • Target table in data lake
  • Warehouse sized for the read workload (larger warehouse = faster read, more cost)

Read Pattern

The Glue Snowflake connector reads via Snowflake's COPY INTO mechanism under the hood -- efficient for large extracts.

snowflake_df = glueContext.create_dynamic_frame.from_options(
    connection_type="snowflake",
    connection_options={
        "connectionName": args['connection_name'],
        "sfDatabase": args['database'],
        "sfSchema": args['schema'],
        "dbtable": args['table']
    }
).toDF()

For custom SQL, use query instead of dbtable:

connection_options={
    "connectionName": args['connection_name'],
    "query": "SELECT id, name, updated_at FROM SALES.ORDERS WHERE status = 'CLOSED'"
}

Incremental Loading

Snowflake has reliable timestamps on most tables. Common watermark columns:

  • Application-maintained updated_at / modified_at
  • Snowflake-maintained _FIVETRAN_SYNCED if sourced via Fivetran
  • INFORMATION_SCHEMA.TABLES.LAST_ALTERED for schema-level freshness (not row-level)

For tables without an updated_at, options:

  • Query SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORY or TABLE_STORAGE_METRICS to identify changed tables for full refresh scheduling
  • Use Snowflake Streams to capture CDC (advanced; requires Snowflake-side setup -- see Snowflake Streams docs)

Standard watermark filter in the custom query:

connection_options={
    "connectionName": args['connection_name'],
    "query": f"SELECT * FROM {source_table} WHERE updated_at > '{last_watermark}'"
}

See incremental-loading.md for watermark storage and the broader incremental pattern.

Partition Pruning

Snowflake tables are automatically micro-partitioned. Push down filters via the query option -- do not pull full tables and filter in Spark.

Clustered tables benefit most from filter push-down. Check cluster keys:

SHOW TABLES LIKE '<table>' IN SCHEMA <db>.<schema>;
-- Look at CLUSTER_BY column

If the source table is clustered on created_date and you filter on created_date >= '2026-01-01', Snowflake prunes micro-partitions and returns only relevant data.

Type Mapping

Snowflake Iceberg Notes
VARCHAR, STRING, TEXT STRING
NUMBER(p,s) DECIMAL(p,s)
NUMBER (no scale) BIGINT
FLOAT, DOUBLE DOUBLE
BOOLEAN BOOLEAN
DATE DATE
TIME STRING Iceberg has no TIME type
TIMESTAMP_NTZ TIMESTAMP Naive timestamp
TIMESTAMP_LTZ, TIMESTAMP_TZ TIMESTAMPTZ Timezone-aware
VARIANT STRING Serialize as JSON
OBJECT STRUCT or STRING Flatten or serialize
ARRAY ARRAY or STRING
BINARY BINARY
GEOGRAPHY, GEOMETRY STRING GeoJSON or WKT

Further Reading

Source: SKILL.md on GitHub

2 warnings17d3 checks · Risk SAFE
  • Gen Agent Trust Hub17d

    This skill facilitates the ingestion of data from various external sources into an AWS data lake. It contains security considerations related to the processing of untrusted data and the dynamic generation of Spark scripts, which are characteristic of ETL (Extract, Transform, Load) operations. These patterns are consistent with the skill's purpose and are used within a managed cloud environment.

  • Socket17d

    1 alert: gptAnomaly

  • Snyk17d

    Risk: MEDIUM · 1 issue

Signed by skilld at b33847d. This ties the file your Agent reads to that commit on GitHub. It does not review the instructions.

Last checked against GitHub yesterday.

Activeupdated 2 months ago
Other metadata
metadata
{
  "version": "1",
  "argument-hint": "'[source-path|connection-name|table-name] [--target s3-tables|iceberg|parquet]'"
}

README badge

README badge for aws/agent-toolkit-for-aws/ingesting-into-data-lake