AWS Athena · 2026-10-01 · Athena arc
AWS Athena for MCP Server Tools: Federated Queries, Iceberg Tables, and Materialized Analytics
Five Athena service areas composed into a production analytics integration for MCP server tools — federated queries via Lambda connectors (register data sources with create_data_catalog(Type="LAMBDA"); query using three-part "catalog"."database"."table" notation; spill bucket configured via spill_bucket Lambda env var — Athena writes intermediate data there when a query result exceeds the Lambda 6 MB response limit; S3 lifecycle rule required to expire spill objects daily or they accumulate; DynamoDB connector pushes equality predicates on the partition key down to KeyConditionExpression — missing the partition key predicate causes a full-table Scan across all shards; custom connectors implement MetadataHandler + RecordHandler from the athena-federation-sdk and expose splits as the unit of parallelism; cross-source JOINs pull the smaller side as a broadcast hash join — keep the federated side narrow to limit Lambda invocations per split), workgroup cost isolation (create workgroups via create_work_group() with EnforceWorkGroupConfiguration: True — this overrides whatever OutputLocation client code passes in start_query_execution, critical when caller code cannot be trusted to route results to the right bucket; BytesScannedCutoffPerQuery cancels runaway queries with a StateChangeReason containing "BytesScannedCutoffPerQuery"; DataScannedInBytes from get_query_execution Statistics is the per-query cost signal; custom MCPServer/AthenaCosts CloudWatch namespace for per-tool attribution; SSE_KMS with a dedicated CMK for regulated workloads; workgroup tags for Cost Explorer chargeback), Iceberg tables for transactional MCP data tools (TBLPROPERTIES('table_type'='ICEBERG') enables row-level UPDATE/DELETE/MERGE without partition rewrite; MERGE INTO for upserts with WHEN MATCHED UPDATE, WHEN MATCHED DELETE, WHEN NOT MATCHED INSERT; time-travel via FOR TIMESTAMP AS OF or FOR VERSION AS OF snapshot_id; virtual metadata tables $snapshots, $history, $manifests, $partitions, $files; schema evolution — ADD/RENAME/DROP COLUMN, CHANGE COLUMN for type widening, ADD PARTITION FIELD with transform functions; GDPR erasure via DELETE + VACUUM; OPTIMIZE BIN_PACK compaction; VACUUM RETAIN N DAYS for snapshot and orphan file cleanup; Athena does not enforce distributed concurrency across multiple writers — design so only one writer path touches a given partition at a time), query result caching ($5/TB Athena pricing makes repeated identical queries expensive — hash the normalized SQL, store the completed QueryExecutionId with TTL in DynamoDB, and serve subsequent calls via GetQueryResults on the cached ID — reading the existing S3 result file incurs no scan charge; GetQueryResults pagination: MaxResults: 1000, NextToken, first row in first page is the column header row; for result sets over 10,000 rows, read the S3 CSV directly via get_object or stream via io.TextIOWrapper; cache TTL must be under 45 days — Athena's default result file retention; proactive invalidation via EventBridge Glue partition events), and CTAS materialization for pre-computed analytics layers (WITH(format='PARQUET', parquet_compression='SNAPPY', external_location, partitioned_by, bucketed_by, bucket_count); CTAS fails if the external location already contains objects — always delete S3 objects and drop the Glue partition before refresh; partition column must be last in the SELECT column list; bucketing effective for JOINs only when both sides use the same column and the same bucket_count; UNLOAD for one-time exports without Glue catalog registration; incremental refresh by listing existing Glue partitions and materializing only the missing date windows). This guide synthesizes the operational mechanics that matter most when MCP server tools drive Athena across the query execution, data storage, and analytics materialization dimensions.
Pattern 1 — Query execution model: federated sources and workgroup cost isolation
MCP server tools that need to query data beyond the S3 data lake have two problems to solve simultaneously: how to reach operational data sources without ETL pipelines, and how to prevent any single tool from generating a surprise Athena bill. Federated queries solve the first problem; workgroups solve the second. Getting both wrong in a multi-tenant MCP server means users trigger expensive full-table scans on production RDS instances without any budget guardrails.
Federated queries — registering and calling Lambda connectors
The federated query architecture uses a Lambda function (the connector) as a translation layer between Athena's SQL engine and the target data source. Deploy the connector once from the Serverless Application Repository, then register it as a named data catalog:
import boto3
import time
athena = boto3.client("athena", region_name="us-east-1")
lambda_client = boto3.client("lambda", region_name="us-east-1")
# Register the connector Lambda as an Athena data source
athena.create_data_catalog(
Name="rds-orders-catalog",
Type="LAMBDA",
Description="Federated access to RDS orders database via JDBC connector",
Parameters={
"function": "arn:aws:lambda:us-east-1:123456789012:function:athena-rds-connector",
},
)
# Configure the connector's spill bucket via Lambda environment variables.
# Athena writes intermediate data here when a query result exceeds the Lambda 6 MB limit.
# Without a spill bucket, large federated queries fail with a serialization error.
lambda_client.update_function_configuration(
FunctionName="athena-rds-connector",
Environment={
"Variables": {
"spill_bucket": "my-athena-spill-bucket",
"spill_prefix": "federation-spill/",
"jdbc_connection_string": "jdbc:postgresql://rds-host:5432/orders",
"secret_manager_enabled": "true",
"secret_name": "rds/orders/credentials",
}
},
)
# Lifecycle rule to expire spill objects — they accumulate quickly during heavy query loads
s3 = boto3.client("s3")
s3.put_bucket_lifecycle_configuration(
Bucket="my-athena-spill-bucket",
LifecycleConfiguration={
"Rules": [{
"ID": "expire-federation-spill",
"Filter": {"Prefix": "federation-spill/"},
"Status": "Enabled",
"Expiration": {"Days": 1},
}]
},
)
Once the catalog is registered, MCP tool handlers query across data sources using the three-part "catalog"."database"."table" naming convention. The S3/Glue default catalog is "AwsDataCatalog"; the federated connector uses the name you chose in create_data_catalog:
def run_cross_source_query(user_id: str) -> list[dict]:
"""MCP tool: join S3 event data with live RDS customer records."""
sql = f"""
SELECT
o.order_id,
o.order_date,
o.total_amount,
c.customer_name,
c.email,
c.tier
FROM "AwsDataCatalog"."analytics"."orders" AS o
JOIN "rds-orders-catalog"."public"."customers" AS c
ON o.customer_id = c.id
WHERE o.user_id = '{user_id}'
AND o.order_date >= DATE '2026-09-01'
AND o.status = 'completed'
ORDER BY o.total_amount DESC
LIMIT 100
"""
start_resp = athena.start_query_execution(
QueryString=sql,
WorkGroup="mcp-reporting-tools",
ResultConfiguration={"OutputLocation": "s3://my-athena-results/"},
)
qid = start_resp["QueryExecutionId"]
while True:
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
state = exec_resp["QueryExecution"]["Status"]["State"]
if state == "SUCCEEDED":
break
elif state in ("FAILED", "CANCELLED"):
reason = exec_resp["QueryExecution"]["Status"].get("StateChangeReason", "")
raise RuntimeError(f"Federated query {state}: {reason}")
time.sleep(2)
rows = []
columns = None
paginator = athena.get_paginator("get_query_results")
for page in paginator.paginate(QueryExecutionId=qid):
result_rows = page["ResultSet"]["Rows"]
if columns is None:
columns = [col["VarCharValue"] for col in result_rows[0]["Data"]]
result_rows = result_rows[1:]
for row in result_rows:
rows.append({columns[i]: cell.get("VarCharValue") for i, cell in enumerate(row["Data"])})
return rows
The critical behavior to understand for DynamoDB federation: the AthenaDynamoDBConnector translates SQL equality predicates on the partition key into KeyConditionExpression calls — a fast, index-bound read. Any query that omits a partition key predicate falls back to a full-table Scan across all shards, which is both expensive and slow. For MCP tools that query DynamoDB via Athena, always ensure the SQL WHERE clause includes a partition key equality condition.
Workgroup cost isolation — enforced budgets per tool
Workgroups assign separate S3 result buckets, scan limits, and CloudWatch metrics to each MCP tool category. The EnforceWorkGroupConfiguration: True flag is the critical setting — without it, client code can override the result location and bypass the cost controls:
def create_mcp_workgroup(
name: str,
result_bucket: str,
result_prefix: str,
bytes_scanned_limit: int,
description: str = "",
) -> None:
athena.create_work_group(
Name=name,
Description=description,
Configuration={
"ResultConfiguration": {
"OutputLocation": f"s3://{result_bucket}/{result_prefix}/",
"EncryptionConfiguration": {
"EncryptionOption": "SSE_S3",
},
},
# EnforceWorkGroupConfiguration: True means the workgroup's OutputLocation
# overrides whatever OutputLocation the client passes in start_query_execution.
"EnforceWorkGroupConfiguration": True,
"PublishCloudWatchMetricsEnabled": True,
"BytesScannedCutoffPerQuery": bytes_scanned_limit,
"EngineVersion": {
"SelectedEngineVersion": "Athena engine version 3",
},
},
Tags=[
{"Key": "mcp-tool", "Value": name},
{"Key": "cost-center", "Value": "analytics"},
],
)
create_mcp_workgroup(
name="mcp-reporting-tools",
result_bucket="my-athena-results",
result_prefix="reporting",
bytes_scanned_limit=50 * 1024**3,
description="Workgroup for MCP reporting and dashboard tools",
)
create_mcp_workgroup(
name="mcp-exploration-tools",
result_bucket="my-athena-results",
result_prefix="exploration",
bytes_scanned_limit=5 * 1024**3,
description="Workgroup for MCP ad-hoc query and exploration tools",
)
create_mcp_workgroup(
name="mcp-etl-jobs",
result_bucket="my-athena-results",
result_prefix="etl",
bytes_scanned_limit=500 * 1024**3,
description="Workgroup for MCP-triggered ETL and CTAS operations",
)
When a query exceeds BytesScannedCutoffPerQuery, Athena cancels it with a StateChangeReason containing "BytesScannedCutoffPerQuery". Surface this to the MCP tool caller with a clear message — it indicates the query needs a partition filter, not a retry:
import cloudwatch_client
def run_mcp_query(sql: str, tool_name: str, workgroup: str) -> tuple[list[dict], dict]:
start_resp = athena.start_query_execution(
QueryString=sql,
WorkGroup=workgroup,
ResultConfiguration={"OutputLocation": "s3://ignored/by-workgroup/"},
)
qid = start_resp["QueryExecutionId"]
while True:
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
state = exec_resp["QueryExecution"]["Status"]["State"]
if state == "SUCCEEDED":
stats = exec_resp["QueryExecution"]["Statistics"]
scanned_bytes = stats.get("DataScannedInBytes", 0)
cost_usd = max(scanned_bytes, 10 * 1024**2) / (1024**4) * 5.0
# Emit per-tool cost metrics to a custom CloudWatch namespace
cloudwatch = boto3.client("cloudwatch", region_name="us-east-1")
cloudwatch.put_metric_data(
Namespace="MCPServer/AthenaCosts",
MetricData=[{
"MetricName": "DataScannedBytes",
"Dimensions": [
{"Name": "ToolName", "Value": tool_name},
{"Name": "WorkGroup", "Value": workgroup},
],
"Value": scanned_bytes,
"Unit": "Bytes",
}],
)
rows = _fetch_all_results(qid)
return rows, {"scanned_bytes": scanned_bytes, "cost_usd": cost_usd}
elif state in ("FAILED", "CANCELLED"):
reason = exec_resp["QueryExecution"]["Status"].get("StateChangeReason", "")
if "BytesScannedCutoffPerQuery" in reason:
raise RuntimeError(
"Query cancelled: exceeded workgroup data scan limit. "
"Add a WHERE clause or partition filter to reduce scanned data."
)
raise RuntimeError(f"Athena query {state}: {reason}")
time.sleep(2)
def _fetch_all_results(qid: str) -> list[dict]:
paginator = athena.get_paginator("get_query_results")
rows = []
columns = None
for page in paginator.paginate(QueryExecutionId=qid):
result_rows = page["ResultSet"]["Rows"]
if columns is None:
columns = [col["VarCharValue"] for col in result_rows[0]["Data"]]
result_rows = result_rows[1:]
for row in result_rows:
rows.append({columns[i]: cell.get("VarCharValue") for i, cell in enumerate(row["Data"])})
return rows
For multi-tenant MCP servers where different tenants drive Athena queries, create one workgroup per tenant tier and tag each with the tenant identifier. AWS Cost Explorer aggregates costs by workgroup tag, providing per-tenant chargeback reports without any custom accounting infrastructure. Update an existing workgroup's scan limit without recreating it using athena.update_work_group(WorkGroup=name, ConfigurationUpdates={...}).
Pattern 2 — Iceberg for transactional MCP data tools
Iceberg tables on Athena give MCP data tools transactional writes, time-travel queries, schema evolution without downtime, and row-level deletes — all on S3, without a running database engine. The difference from Hive-style external tables is the metadata layer: every INSERT, UPDATE, DELETE, or MERGE creates an atomic snapshot commit. A load_data tool writing new rows does not break concurrent query_data tool reads. A correct_record tool can issue a targeted UPDATE without rewriting entire partitions. A get_audit_snapshot tool can time-travel to any historical point in the table's commit history.
Iceberg DDL and transactional writes
Create an Iceberg table by adding TBLPROPERTIES('table_type'='ICEBERG') to the DDL. All subsequent DML uses Athena engine version 3 — confirm the workgroup's SelectedEngineVersion is set accordingly:
import boto3
import time
athena = boto3.client("athena", region_name="us-east-1")
def run_ddl(sql: str, workgroup: str = "mcp-etl-jobs") -> str:
resp = athena.start_query_execution(QueryString=sql, WorkGroup=workgroup)
qid = resp["QueryExecutionId"]
while True:
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
state = exec_resp["QueryExecution"]["Status"]["State"]
if state == "SUCCEEDED":
return qid
elif state in ("FAILED", "CANCELLED"):
raise RuntimeError(
f"DDL failed: {exec_resp['QueryExecution']['Status'].get('StateChangeReason')}"
)
time.sleep(2)
# Create Iceberg table — TBLPROPERTIES triggers Iceberg metadata management
run_ddl("""
CREATE TABLE analytics.mcp_events (
event_id STRING,
session_id STRING,
tool_name STRING,
user_id STRING,
event_ts TIMESTAMP,
tool_payload STRING,
status STRING,
latency_ms BIGINT
)
LOCATION 's3://my-datalake/iceberg/mcp_events/'
TBLPROPERTIES (
'table_type' = 'ICEBERG',
'format' = 'parquet',
'write_compression' = 'snappy',
'optimize_rewrite_delete_file_threshold' = '10'
)
""")
# Row-level update — Iceberg rewrites only the affected data files
run_ddl("""
UPDATE analytics.mcp_events
SET status = 'retried', latency_ms = 320
WHERE event_id = 'evt-001'
""")
# Row-level delete — no partition rewrite required
run_ddl("""
DELETE FROM analytics.mcp_events
WHERE user_id = 'user-123' AND event_ts < TIMESTAMP '2026-09-01 00:00:00'
""")
Concurrent INSERT operations on non-overlapping partitions are safe. Concurrent UPDATE or DELETE operations on the same row may produce optimistic-concurrency failures — Athena does not enforce distributed write serialization across multiple simultaneous writers. Design MCP tool architectures so that only one writer path modifies a given partition at a time, typically by routing all writes through a single worker process or a serialized queue.
MERGE INTO for upserts and synchronization
MCP tools that synchronize rows from a staging table or replicate changes from an operational database use MERGE INTO. It is cheaper than a DELETE+INSERT cycle because Iceberg tracks which files were touched and rewrites only those:
run_ddl("""
MERGE INTO analytics.mcp_events AS target
USING (
SELECT * FROM analytics.mcp_events_staging
) AS source
ON target.event_id = source.event_id
WHEN MATCHED AND source.status = 'deleted' THEN
DELETE
WHEN MATCHED THEN
UPDATE SET
status = source.status,
latency_ms = source.latency_ms,
tool_payload = source.tool_payload
WHEN NOT MATCHED THEN
INSERT (event_id, session_id, tool_name, user_id, event_ts, tool_payload, status, latency_ms)
VALUES (source.event_id, source.session_id, source.tool_name, source.user_id,
source.event_ts, source.tool_payload, source.status, source.latency_ms)
""")
Time-travel queries for audit and GDPR erasure
Iceberg snapshots let MCP audit tools query the table as it existed at any historical point. Two forms are available — timestamp-based (more human-readable) and snapshot-ID-based (immune to clock skew and exact for reproducibility):
# Time-travel by timestamp — query data as it existed at a specific moment
run_ddl("""
SELECT * FROM analytics.mcp_events
FOR TIMESTAMP AS OF TIMESTAMP '2026-09-15 00:00:00'
WHERE tool_name = 'query_data'
""")
# List available snapshots from the virtual metadata table
def list_snapshots() -> list[dict]:
qid = run_ddl("""
SELECT snapshot_id, committed_at, operation, summary
FROM "analytics"."mcp_events$snapshots"
ORDER BY committed_at DESC
LIMIT 20
""")
return _fetch_all_results(qid)
# Time-travel by snapshot ID — deterministic for reproducible reports
snapshot_id = "5789234567890123456"
run_ddl(f"""
SELECT count(*) AS event_count, tool_name
FROM analytics.mcp_events
FOR VERSION AS OF {snapshot_id}
GROUP BY tool_name
""")
# $history shows snapshot transitions — useful for debugging UPDATE/DELETE chains
# Each row: snapshot_id, parent_id, is_current_ancestor, made_current_at
qid = run_ddl("""
SELECT * FROM "analytics"."mcp_events$history"
ORDER BY made_current_at DESC
LIMIT 10
""")
GDPR "right to be forgotten" erasure via Iceberg: issue a targeted DELETE FROM for the user's rows, then run VACUUM to expire old snapshots that still contain the deleted data. After VACUUM, the deleted rows are no longer accessible via time-travel — the snapshot containing them has been physically removed from the S3 metadata files. Set the VACUUM retention period to match your time-travel depth requirement (the minimum window you need for auditing), not shorter.
Schema evolution and table maintenance
Iceberg tracks schema versions in metadata, so adding or renaming columns does not require rewriting data files. Existing Parquet files return NULL for newly added columns transparently:
# Add columns — existing data files return NULL for the new columns
run_ddl("""
ALTER TABLE analytics.mcp_events
ADD COLUMNS (error_code STRING, retry_count INT)
""")
# Rename — Iceberg maps the old column ID to the new name in metadata
# No data files are rewritten
run_ddl("""
ALTER TABLE analytics.mcp_events
RENAME COLUMN latency_ms TO latency_ms_p50
""")
# Type widening (INT → BIGINT, FLOAT → DOUBLE) — only widening is supported
run_ddl("""
ALTER TABLE analytics.mcp_events
CHANGE COLUMN retry_count retry_count BIGINT
""")
# Add a new partition field — old data retains the old partition scheme;
# new data uses the new one; Athena reads both transparently
run_ddl("""
ALTER TABLE analytics.mcp_events
ADD PARTITION FIELD day(event_ts)
""")
# Compact small files after heavy INSERT/UPDATE/DELETE activity.
# OPTIMIZE rewrites files in the specified window; old files become orphans until VACUUM.
run_ddl("""
OPTIMIZE analytics.mcp_events
REWRITE DATA USING BIN_PACK
WHERE event_ts >= TIMESTAMP '2026-10-01 00:00:00'
""")
# Remove expired snapshots and orphan files.
# After OPTIMIZE, old file versions are unreachable but still stored on S3 until VACUUM.
run_ddl("""
VACUUM analytics.mcp_events
RETAIN 7 DAYS
""")
Iceberg's hidden partitioning means partition columns are never stored in data files as explicit columns — they are derived from source column values using transform functions (day, month, bucket, truncate). MCP tool SQL never needs explicit partition filters; Iceberg's metadata pruning applies partition elimination automatically based on WHERE predicates on the underlying source column. Schedule OPTIMIZE as a periodic MCP admin tool operation to prevent small-file accumulation from degrading scan performance on heavily-written tables.
Pattern 3 — Materialized analytics layer: CTAS and query result caching
MCP analytics tools that serve reporting queries face the same challenge: Athena charges $5/TB scanned, so running the same expensive aggregation for every tool invocation is both slow and costly. The solution is a two-layer materialization approach — CTAS to pre-aggregate raw events into compact Parquet tables, and SQL-hash caching to reuse completed query executions without rescanning data. Together they can cut MCP tool query costs by 50–100× for workloads with any repetition.
Partitioned CTAS for daily materialized aggregations
CTAS writes the result of a SELECT into a new external table in the Glue Data Catalog. The critical failure mode: CTAS fails if the external_location already contains objects. Always delete the S3 objects and drop the Glue partition before refreshing a given date partition:
import boto3
import time
from datetime import date, timedelta
athena = boto3.client("athena", region_name="us-east-1")
def materialize_daily_summary(target_date: date) -> None:
"""Materialize MCP event summary for a given date partition."""
dt_str = target_date.isoformat()
s3_prefix = f"s3://my-datalake/materialized/mcp_daily_summary/dt={dt_str}/"
bucket, prefix = "my-datalake", f"materialized/mcp_daily_summary/dt={dt_str}/"
s3 = boto3.client("s3")
paginator = s3.get_paginator("list_objects_v2")
# Step 1: Delete existing objects at this partition location
# CTAS fails with an error if external_location is non-empty
for page in paginator.paginate(Bucket=bucket, Prefix=prefix):
objects = page.get("Contents", [])
if objects:
s3.delete_objects(
Bucket=bucket,
Delete={"Objects": [{"Key": obj["Key"]} for obj in objects]},
)
# Step 2: Drop old Glue partition
glue = boto3.client("glue", region_name="us-east-1")
try:
glue.delete_partition(
DatabaseName="analytics",
TableName="mcp_daily_summary",
PartitionValues=[dt_str],
)
except glue.exceptions.EntityNotFoundException:
pass
# Step 3: CTAS — partition column must be last in the SELECT list
resp = athena.start_query_execution(
QueryString=f"""
CREATE TABLE analytics.mcp_daily_summary_temp_{dt_str.replace('-', '')}
WITH (
format = 'PARQUET',
parquet_compression = 'SNAPPY',
external_location = '{s3_prefix}',
partitioned_by = ARRAY['dt']
)
AS
SELECT
tool_name,
status,
COUNT(*) AS call_count,
COUNT(*) FILTER (WHERE status = 'success') AS success_count,
COUNT(*) FILTER (WHERE status != 'success') AS error_count,
AVG(latency_ms) AS avg_latency_ms,
APPROX_PERCENTILE(latency_ms, 0.95) AS p95_latency_ms,
APPROX_PERCENTILE(latency_ms, 0.99) AS p99_latency_ms,
'{dt_str}' AS dt
FROM analytics.mcp_events
WHERE event_ts >= TIMESTAMP '{dt_str} 00:00:00'
AND event_ts < TIMESTAMP '{dt_str} 00:00:00' + INTERVAL '1' DAY
GROUP BY tool_name, status
""",
WorkGroup="mcp-etl-jobs",
)
qid = resp["QueryExecutionId"]
while True:
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
state = exec_resp["QueryExecution"]["Status"]["State"]
if state == "SUCCEEDED":
stats = exec_resp["QueryExecution"]["Statistics"]
print(f"CTAS {dt_str}: {stats.get('DataScannedInBytes', 0) / 1024**3:.2f} GB scanned")
break
elif state in ("FAILED", "CANCELLED"):
raise RuntimeError(
f"CTAS failed: {exec_resp['QueryExecution']['Status'].get('StateChangeReason')}"
)
time.sleep(5)
# Step 4: Register the new partition in the Glue catalog
glue.create_partition(
DatabaseName="analytics",
TableName="mcp_daily_summary",
PartitionInput={
"Values": [dt_str],
"StorageDescriptor": {
"Location": s3_prefix,
"InputFormat": "org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat",
"OutputFormat": "org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat",
"SerdeInfo": {
"SerializationLibrary": "org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe"
},
},
},
)
Incremental partition refresh keeps CTAS costs proportional to new data volume — a two-year table is no more expensive to refresh daily than a one-week table, because each refresh only touches the new partition's source data:
def refresh_missing_partitions(start_date: date, end_date: date) -> list[str]:
"""Materialize only the date partitions not yet in the summary table."""
glue = boto3.client("glue", region_name="us-east-1")
paginator = glue.get_paginator("get_partitions")
existing = set()
for page in paginator.paginate(DatabaseName="analytics", TableName="mcp_daily_summary"):
for p in page["Partitions"]:
existing.add(p["Values"][0])
missing = []
current = start_date
while current <= end_date:
if current.isoformat() not in existing:
missing.append(current)
current += timedelta(days=1)
for dt in missing:
materialize_daily_summary(dt)
return [d.isoformat() for d in missing]
# Called from an MCP admin tool or scheduled Lambda
refresh_missing_partitions(
start_date=date(2026, 9, 1),
end_date=date.today() - timedelta(days=1),
)
SQL-hash caching to eliminate redundant Athena scans
Even against the materialized summary table, MCP tools that run the same aggregation query repeatedly on the same partition pay full scan cost on each call. SQL-hash caching eliminates this: hash the normalized SQL, store the completed QueryExecutionId with a TTL, and serve subsequent identical calls via GetQueryResults on the cached ID. Reading from an existing S3 result file costs nothing — no rescan, no charge:
import hashlib
import re
from typing import Optional
dynamodb = boto3.resource("dynamodb", region_name="us-east-1")
cache_table = dynamodb.Table("athena-query-cache")
def normalize_sql(sql: str) -> str:
sql = re.sub(r'--[^\n]*', '', sql)
sql = re.sub(r'/\*.*?\*/', '', sql, flags=re.DOTALL)
sql = re.sub(r'\s+', ' ', sql).strip().lower()
return sql
def get_cached_qid(sql: str) -> Optional[str]:
sql_hash = hashlib.sha256(normalize_sql(sql).encode()).hexdigest()
item = cache_table.get_item(Key={"sql_hash": sql_hash}).get("Item")
if not item:
return None
if int(item.get("expires_at", 0)) < int(time.time()):
return None # expired
return item["query_execution_id"]
def store_cached_qid(sql: str, qid: str, ttl_seconds: int) -> None:
sql_hash = hashlib.sha256(normalize_sql(sql).encode()).hexdigest()
cache_table.put_item(Item={
"sql_hash": sql_hash,
"query_execution_id": qid,
"sql_preview": sql[:200],
"cached_at": int(time.time()),
"expires_at": int(time.time()) + ttl_seconds,
})
# TTL by freshness requirement — must be under 45 days (Athena result file retention default)
CACHE_TTL_BY_TOOL = {
"get_realtime_metrics": 15 * 60, # 15 min — near-real-time dashboards
"get_daily_report": 60 * 60, # 1 hr — daily aggregation is stable
"get_historical_trend": 24 * 60 * 60, # 24 hr — past periods never change
"list_available_tables": 4 * 60 * 60, # 4 hr — schema changes slowly
"run_custom_query": 0, # 0 = always execute fresh
}
def run_mcp_tool_query(tool_name: str, sql: str, workgroup: str) -> list[dict]:
ttl = CACHE_TTL_BY_TOOL.get(tool_name, 3600)
if ttl == 0:
resp = athena.start_query_execution(QueryString=sql, WorkGroup=workgroup)
_wait_for_completion(resp["QueryExecutionId"])
return _fetch_all_results(resp["QueryExecutionId"])
cached_qid = get_cached_qid(sql)
if cached_qid:
try:
exec_resp = athena.get_query_execution(QueryExecutionId=cached_qid)
if exec_resp["QueryExecution"]["Status"]["State"] == "SUCCEEDED":
return _fetch_all_results(cached_qid)
except athena.exceptions.InvalidRequestException:
pass # execution ID expired — fall through to re-run
resp = athena.start_query_execution(QueryString=sql, WorkGroup=workgroup)
qid = resp["QueryExecutionId"]
_wait_for_completion(qid)
store_cached_qid(sql, qid, ttl)
return _fetch_all_results(qid)
def _wait_for_completion(qid: str, timeout: int = 300) -> None:
deadline = time.time() + timeout
while time.time() < deadline:
resp = athena.get_query_execution(QueryExecutionId=qid)
state = resp["QueryExecution"]["Status"]["State"]
if state == "SUCCEEDED":
return
elif state in ("FAILED", "CANCELLED"):
raise RuntimeError(
f"Query {state}: {resp['QueryExecution']['Status'].get('StateChangeReason')}"
)
time.sleep(2)
raise TimeoutError(f"Athena query timed out after {timeout}s")
For result sets larger than roughly 10,000 rows, paginating through GetQueryResults at 1,000 rows per page becomes the latency bottleneck — each page is a separate HTTP round-trip. Read the S3 result CSV directly instead:
import csv
import io
s3 = boto3.client("s3")
def fetch_large_results_from_s3(qid: str) -> list[dict]:
"""Read Athena result CSV directly from S3 — faster than GetQueryResults for large sets."""
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
output_location = exec_resp["QueryExecution"]["ResultConfiguration"]["OutputLocation"]
path = output_location.replace("s3://", "")
bucket, key = path.split("/", 1)
response = s3.get_object(Bucket=bucket, Key=key)
content = response["Body"].read().decode("utf-8")
return list(csv.DictReader(io.StringIO(content)))
def stream_large_results_from_s3(qid: str, chunk_size: int = 1000):
"""Stream result CSV from S3 in row chunks to avoid loading the entire file into memory."""
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
output_location = exec_resp["QueryExecution"]["ResultConfiguration"]["OutputLocation"]
path = output_location.replace("s3://", "")
bucket, key = path.split("/", 1)
response = s3.get_object(Bucket=bucket, Key=key)
reader = csv.DictReader(io.TextIOWrapper(response["Body"], encoding="utf-8"))
chunk = []
for row in reader:
chunk.append(row)
if len(chunk) >= chunk_size:
yield chunk
chunk = []
if chunk:
yield chunk
Cache invalidation should happen proactively when source data changes. Subscribe to Glue catalog events via EventBridge — whenever a new partition is registered (a new day's data finishes loading), fire a Lambda that calls cache_table.delete_item for the SQL hashes that read from the updated table. This prevents stale results from appearing in MCP tool responses even when the TTL hasn't expired yet.
Consolidated failure modes across the Athena arc
| Surface | Failure mode | Root cause | Correct response |
|---|---|---|---|
| Federated queries | Connector spill errors on large queries | spill_bucket environment variable not configured on the connector Lambda; Athena cannot write intermediate data when the result exceeds the Lambda 6 MB response limit |
Set spill_bucket and spill_prefix on the connector Lambda; add S3 lifecycle rule to expire spill objects after 1 day |
| Federated queries | Full DynamoDB table scan instead of index lookup | SQL predicate does not include a partition key equality condition; the DynamoDB connector falls back to a full-table Scan API call | Always include a partition key equality predicate (WHERE user_id = '...') when querying DynamoDB via Athena federation; test query plans before deploying |
| Workgroups | Client code routes results to wrong S3 bucket | EnforceWorkGroupConfiguration not set to True; client-supplied OutputLocation overrides the workgroup's configured result location |
Set EnforceWorkGroupConfiguration: True on every workgroup; workgroup configuration is then authoritative regardless of what client code passes |
| Workgroups | Query silently cancelled with no clear error | Query exceeded BytesScannedCutoffPerQuery; Athena cancels with StateChangeReason containing "BytesScannedCutoffPerQuery" — MCP tool error handling does not surface this reason |
Check StateChangeReason explicitly on CANCELLED queries and surface a "add a partition filter" message rather than a generic query failure |
| Iceberg | CTAS fails on partition refresh | CTAS fails if the external_location already contains objects from a prior run; S3 objects and Glue partition not cleaned up before re-running |
Always delete S3 objects under the target prefix and call glue.delete_partition() before running the CTAS refresh; wrap in try/except for EntityNotFoundException on first run |
| Iceberg | Concurrent UPDATE/DELETE conflicts | Multiple MCP server instances write to the same Iceberg table partition simultaneously; Athena does not enforce distributed write serialization across independent processes | Route all writes through a single worker process per partition or use a serialized queue; concurrent INSERTs on non-overlapping partitions are safe |
| Iceberg | GDPR deletion remains accessible via time-travel | DELETE statement removes rows from current table view but old snapshots containing deleted rows remain on S3 until VACUUM is run | Issue VACUUM ... RETAIN N DAYS after GDPR deletions; set retention period to the minimum time-travel depth required, not longer |
| Iceberg | Query performance degrades over time | Heavy INSERT/UPDATE/DELETE activity produces many small data files per partition; Athena S3 GET count per query grows until latency becomes unacceptable | Schedule OPTIMIZE ... REWRITE DATA USING BIN_PACK as a periodic maintenance operation; follow with VACUUM to remove orphan files |
| Query caching | Cached execution ID no longer valid | Athena result files are retained for 45 days by default; a cached QueryExecutionId with a TTL longer than 45 days causes InvalidRequestException on GetQueryResults |
Set all cache TTLs under 45 days; validate the cached execution ID with get_query_execution before calling GetQueryResults and fall through to a fresh run on InvalidRequestException |
| Query caching | Stale cache results after ETL partition load | Cache TTL has not expired but the underlying data changed; a new partition was loaded or an existing partition was refreshed | Subscribe to Glue catalog events via EventBridge; proactively invalidate cache entries for affected SQL patterns when a new partition is registered |
| CTAS | Bucketing does not reduce JOIN shuffle | CTAS bucketing is only effective when both sides of a JOIN are bucketed on the same column with the same bucket_count; asymmetric bucketing causes the engine to fall back to a regular shuffle join |
Bucket both the fact and dimension tables on the same column with the same bucket_count; document the bucketing contract so future schema changes preserve it |
| CTAS | Glue catalog polluted with stale temp tables | CTAS creates a new Glue table on every run; temp tables from intermediate CTAS steps are never dropped, accumulating in the catalog | Call glue.delete_table() to clean up CTAS temp tables after each refresh; data files remain on S3 but the Glue entry is removed |
Production checklists
Federated query checklist
spill_bucketandspill_prefixconfigured on every connector Lambda- S3 lifecycle rule expiring spill objects after 1 day
- DynamoDB connector queries always include partition key equality predicate
- Cross-source JOIN keeps the federated (operational) side narrow with predicate pushdown
- Connector Lambda has a dedicated security group with VPC access to the target database
- Materialized S3 snapshot (daily CTAS) available as fallback for MCP tools that can tolerate slight staleness
Workgroup checklist
EnforceWorkGroupConfiguration: Trueon every workgroupBytesScannedCutoffPerQueryset per workgroup based on tool scan profile — not a single global limitPublishCloudWatchMetricsEnabled: Trueon all workgroups- Custom
MCPServer/AthenaCostsnamespace metric emitted withToolNamedimension after every successful query - CloudWatch alarm on
BytesScannedPerQuerydistribution to detect runaway queries before billing period ends - Workgroup tags set for Cost Explorer chargeback reporting
- SSE_KMS with dedicated CMK for workgroups handling regulated data
Iceberg checklist
TBLPROPERTIES('table_type'='ICEBERG')in CREATE TABLE; workgroup'sSelectedEngineVersionis Athena engine version 3- Single writer path per partition — no concurrent UPDATE/DELETE from multiple processes
- GDPR DELETE followed by
VACUUM RETAIN N DAYSbefore confirming erasure completion - VACUUM retention period set to minimum time-travel depth required — not longer
OPTIMIZE ... REWRITE DATA USING BIN_PACKscheduled as periodic maintenance after heavy write periods- MERGE INTO used for upserts instead of DELETE+INSERT — cheaper and atomic
- Schema evolution (ADD/RENAME/DROP COLUMN) tested in non-production before applying to production tables
Query caching and CTAS checklist
- All cache TTLs under 45 days (Athena result file retention default)
- Cached
QueryExecutionIdvalidated withget_query_executionbefore callingGetQueryResults - Cache invalidation triggered by EventBridge Glue partition events — not only by TTL expiry
- Large result sets (>10,000 rows) read from S3 directly via
get_objectrather than paginatingGetQueryResults - CTAS refresh: S3 objects deleted + Glue partition dropped before each run — CTAS fails on non-empty location
- Partition column is the last column in every CTAS SELECT list
- CTAS bucketing: both sides of JOIN use same column and same
bucket_count - CTAS temp tables dropped via
glue.delete_table()after each refresh to prevent catalog clutter - Incremental partition refresh implemented — full-table CTAS re-run not required for tables with long history
Monitor every Athena-backed MCP endpoint
Athena-backed MCP tools fail in non-obvious ways: a connector Lambda cold-starts and times out, a CTAS job silently produces stale data when the S3 cleanup step fails, a workgroup scan limit cancels queries that just grew past the threshold as data volume increases. AliveMCP probes every MCP endpoint every 60 seconds — alerting your team the moment an Athena-backed analytics tool starts returning failures or degraded responses, before users encounter broken dashboards or missing data.
Join the waitlist →