Guide · AWS Athena · Query Caching · MCP Tools
Athena Query Result Caching for MCP Tools
Athena charges $5 per terabyte scanned — MCP tools that run the same query repeatedly pay the full scan cost every time unless the server implements result caching. Athena itself provides no built-in query result reuse across API calls; caching must be implemented in the MCP server layer. The strategy is: hash the normalized SQL string, store the completed QueryExecutionId against that hash with a TTL, and serve subsequent identical requests from GetQueryResults on the cached execution ID. GetQueryResults on a completed execution is free — it reads from the already-written S3 result file without re-scanning the data lake. Three patterns cover most MCP use cases: SQL-hash caching for deterministic queries run multiple times per session; S3 result file direct reads for large result sets where GetQueryResults pagination is too slow; and TTL-aware invalidation for tools that must reflect data freshness windows.
TL;DR
Hash the normalized SQL string. Look up the hash in a cache (DynamoDB or in-process dict). If found and not expired, call GetQueryResults(QueryExecutionId=cached_id) — this reads from the existing S3 result file, no rescan, no charge. If not found or expired, run the query normally and store the new QueryExecutionId with a TTL. For large result sets, read the S3 CSV directly via presigned URL instead of paginating through GetQueryResults.
SQL-hash caching layer for MCP analytics tools
The cache key is a SHA-256 hash of the normalized SQL. Normalization strips comments, normalizes whitespace, and lowercases identifiers — so semantically identical queries with different formatting share the same cache entry:
import boto3
import hashlib
import json
import re
import time
from typing import Optional
athena = boto3.client("athena", region_name="us-east-1")
dynamodb = boto3.resource("dynamodb", region_name="us-east-1")
cache_table = dynamodb.Table("athena-query-cache")
def normalize_sql(sql: str) -> str:
"""Normalize SQL for stable cache key generation."""
# Remove single-line comments
sql = re.sub(r'--[^\n]*', '', sql)
# Remove multi-line comments
sql = re.sub(r'/\*.*?\*/', '', sql, flags=re.DOTALL)
# Normalize whitespace
sql = re.sub(r'\s+', ' ', sql).strip().lower()
return sql
def get_cached_execution_id(sql: str, ttl_seconds: int) -> Optional[str]:
"""Return a still-valid cached QueryExecutionId, or None."""
sql_hash = hashlib.sha256(normalize_sql(sql).encode()).hexdigest()
response = cache_table.get_item(Key={"sql_hash": sql_hash})
item = response.get("Item")
if not item:
return None
# Check TTL (DynamoDB TTL attribute named "expires_at" in epoch seconds)
if int(item.get("expires_at", 0)) < int(time.time()):
return None # expired — re-run the query
return item["query_execution_id"]
def store_cached_execution_id(
sql: str, query_execution_id: str, ttl_seconds: int
) -> None:
"""Store QueryExecutionId in the cache with TTL."""
sql_hash = hashlib.sha256(normalize_sql(sql).encode()).hexdigest()
cache_table.put_item(Item={
"sql_hash": sql_hash,
"query_execution_id": query_execution_id,
"sql_preview": sql[:200], # for debugging — not used as cache key
"cached_at": int(time.time()),
"expires_at": int(time.time()) + ttl_seconds, # DynamoDB TTL field
})
def run_with_sql_cache(
sql: str,
workgroup: str = "mcp-reporting-tools",
ttl_seconds: int = 3600,
) -> list[dict]:
"""Run Athena query with SQL-hash caching; return result rows."""
# Check cache first
cached_qid = get_cached_execution_id(sql, ttl_seconds)
if cached_qid:
# Verify the cached execution is still accessible (results expire after 45 days)
try:
exec_resp = athena.get_query_execution(QueryExecutionId=cached_qid)
if exec_resp["QueryExecution"]["Status"]["State"] == "SUCCEEDED":
return fetch_query_results(cached_qid)
except athena.exceptions.InvalidRequestException:
pass # execution ID no longer valid — fall through to re-run
# Cache miss: execute fresh query
start_resp = athena.start_query_execution(
QueryString=sql,
WorkGroup=workgroup,
)
qid = start_resp["QueryExecutionId"]
_wait_for_completion(qid)
# Store in cache for subsequent calls
store_cached_execution_id(sql, qid, ttl_seconds)
return fetch_query_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("Athena query timed out")
Paginating large result sets with GetQueryResults
GetQueryResults returns up to 1,000 rows per page. MCP tools that return large result sets must paginate using NextToken:
def fetch_query_results(qid: str, max_rows: Optional[int] = None) -> list[dict]:
"""Paginate through all Athena query results; optional row limit."""
rows = []
columns = None
next_token = None
while True:
kwargs = {"QueryExecutionId": qid, "MaxResults": 1000}
if next_token:
kwargs["NextToken"] = next_token
page = athena.get_query_results(**kwargs)
result_rows = page["ResultSet"]["Rows"]
if columns is None:
# First row in first page is always the column header row
columns = [
col.get("VarCharValue", f"col_{i}")
for i, col in enumerate(result_rows[0]["Data"])
]
result_rows = result_rows[1:]
for row in result_rows:
rows.append({
columns[i]: cell.get("VarCharValue") # VarCharValue is None → column is NULL
for i, cell in enumerate(row["Data"])
})
if max_rows and len(rows) >= max_rows:
return rows
next_token = page.get("NextToken")
if not next_token:
break # all pages consumed
return rows
# Example: MCP tool with result size cap to avoid overloading agent context
results = fetch_query_results(qid, max_rows=500)
For result sets larger than ~10,000 rows, GetQueryResults pagination becomes slow — each page is a separate HTTP round-trip. At that scale, read the S3 result CSV directly:
import csv
import io
import boto3
s3 = boto3.client("s3")
def fetch_results_from_s3(qid: str) -> list[dict]:
"""Read Athena result CSV from S3 directly — faster than GetQueryResults for large sets."""
# GetQueryExecution tells us where the result file is
exec_resp = athena.get_query_execution(QueryExecutionId=qid)
output_location = exec_resp["QueryExecution"]["ResultConfiguration"]["OutputLocation"]
# Parse s3://bucket/prefix/qid.csv
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")
reader = csv.DictReader(io.StringIO(content))
return list(reader)
# For very large results, stream in chunks to avoid loading entire CSV into memory
def stream_results_from_s3(qid: str, chunk_size: int = 1000):
"""Stream Athena result CSV from S3 in row chunks."""
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
TTL strategy for different MCP tool freshness requirements
The TTL on a cached Athena result should match the acceptable data staleness for that tool's use case:
# Athena query result files are stored for 45 days by default
# (configurable via workgroup or individual query result bucket lifecycle policy)
# Cache TTL must be shorter than 45 days or the cached execution ID becomes invalid
CACHE_TTL_BY_TOOL = {
# Metrics dashboards: acceptable to show data 15 minutes old
"get_realtime_metrics": 15 * 60,
# Daily reports: acceptable to show data up to 1 hour old
"get_daily_report": 60 * 60,
# Historical trend: data doesn't change for past periods — 24h TTL safe
"get_historical_trend": 24 * 60 * 60,
# Schema/catalog info: changes rarely — 4h TTL
"list_available_tables": 4 * 60 * 60,
# Live query: no cache — always fresh
"run_custom_query": 0, # TTL=0 means 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:
# No cache — run fresh every time
resp = athena.start_query_execution(QueryString=sql, WorkGroup=workgroup)
_wait_for_completion(resp["QueryExecutionId"])
return fetch_query_results(resp["QueryExecutionId"])
return run_with_sql_cache(sql, workgroup=workgroup, ttl_seconds=ttl)
# Cache invalidation: force-invalidate when source data is known to have changed
def invalidate_cache(sql: str) -> None:
sql_hash = hashlib.sha256(normalize_sql(sql).encode()).hexdigest()
cache_table.delete_item(Key={"sql_hash": sql_hash})
When the underlying data changes (a new ETL job completes, a new partition is added to the data lake), proactively call invalidate_cache for the SQL patterns that read from the updated tables. This prevents stale results from appearing in MCP tool responses even when the TTL has not yet expired. For Glue-partitioned tables, subscribe to Glue catalog events via EventBridge to trigger cache invalidation whenever a new partition is registered.
Monitor Athena-backed MCP tools with AliveMCP
Caching infrastructure (DynamoDB cache table, S3 result bucket lifecycle policies) can fail silently — a misconfigured TTL or expired execution ID causes every MCP tool call to fall back to full rescans, tripling your Athena bill. AliveMCP watches every MCP endpoint every 60 seconds so you know when query performance degrades.
Join the waitlist →