Add analytics dataset: projection table + Athena workgroup (Phase 3)

CDK-define the apm_wo_snapshots Glue table with partition projection over
analytics/ (dt as projected date partition, 2026-01-01..NOW). Projection means
no crawler, no MSCK REPAIR, and — critically — the classifier needs no Glue
catalog access at all.

pipeline_stack.py: glue.CfnTable (Parquet SerDe, 17-column schema mirroring the
classifier's snapshot incl. contractor_description and the two boolean flags) +
athena.CfnWorkGroup `apm-wo-analysis` (enforced result location, SSE-S3) + a
30-day lifecycle rule on athena-results/ (disposable query output in a RETAIN
bucket). Trim the classifier role: drop the entire Glue policy statement.

handler.py: stop registering the table at runtime — drop database=/table= from
to_parquet so the classifier writes pure Parquet; partition projection handles
the rest. Keeps overwrite_partitions for idempotent same-day re-uploads.

test_pipeline_synth.py: offline synth assertions (bundling skipped) — projection
properties, column schema/types, workgroup result enforcement, and that no IAM
policy grants glue:* to the classifier.

cdk synth green; 11/11 tests pass.
This commit is contained in:
Adam Moussa 2026-05-28 17:24:36 -04:00
parent 0272db1716
commit 6ee218eea3
3 changed files with 229 additions and 35 deletions

View file

@ -1,8 +1,8 @@
"""Pipeline stack: S3, classifier + slack-post Lambdas, Glue, Athena, IAM.
Phase 1 — exports bucket + drop-folder uploader IAM user.
Phase 2 — Glue database + S3-triggered classifier Lambda (this file).
Phase 3 — Athena workgroup + partition-projection table (TODO).
Phase 2 — Glue database + S3-triggered classifier Lambda.
Phase 3 — partition-projection Glue table + Athena workgroup (this file).
Phase 4 — slack-post Lambda + scoped IAM (TODO).
Lambdas: Python 3.12, ARM64, explicit LogGroup with 60-day retention.
@ -17,6 +17,9 @@ from aws_cdk import (
RemovalPolicy,
Stack,
)
from aws_cdk import (
aws_athena as athena,
)
from aws_cdk import (
aws_glue as glue,
)
@ -42,9 +45,37 @@ from constructs import Construct
GLUE_DATABASE = "apm_wo_analysis"
GLUE_TABLE = "apm_wo_snapshots"
ATHENA_WORKGROUP = "apm-wo-analysis"
ANTHROPIC_SECRET = "apm-wo-analysis/anthropic-api-key"
LAMBDAS_DIR = os.path.join(os.path.dirname(__file__), "..", "..", "lambdas")
# Snapshot schema — mirrors the per-WO record written by the classifier
# (lambdas/classifier/handler.py _build_snapshot). Order/names must match the
# Parquet columns; `dt` is the projected partition key, not a stored column.
SNAPSHOT_COLUMNS = [
("wo_number", "string"),
("wo_description", "string"),
("equipment_code", "string"),
("site", "string"),
("due_date", "string"),
("department", "string"),
("wo_status", "string"),
("hold_reason", "string"),
("last_comment", "string"),
("last_comment_by", "string"),
("last_comment_date", "string"),
("contractor", "string"),
("contractor_description", "string"),
("category", "string"),
("is_escalation", "boolean"),
("is_action", "boolean"),
("mismatch", "string"),
]
_PARQUET_INPUT = "org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat"
_PARQUET_OUTPUT = "org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat"
_PARQUET_SERDE = "org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe"
class PipelineStack(Stack):
def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None:
@ -65,7 +96,14 @@ class PipelineStack(Stack):
id="expire-raw-exports",
prefix="raw/",
expiration=Duration.days(90),
)
),
# Athena query output is disposable; don't let it accumulate in
# a RETAIN bucket. analytics/ snapshots are kept indefinitely.
s3.LifecycleRule(
id="expire-athena-results",
prefix="athena-results/",
expiration=Duration.days(30),
),
],
)
@ -84,9 +122,7 @@ class PipelineStack(Stack):
)
)
# Phase 2/3 — Glue database for the analytics dataset. The classifier
# registers the apm_wo_snapshots table/partitions into it via awswrangler;
# Phase 3 swaps in the partition-projection table definition + Athena.
# Phase 2/3 — Glue database for the analytics dataset.
glue.CfnDatabase(
self,
"AnalyticsDb",
@ -94,6 +130,63 @@ class PipelineStack(Stack):
database_input=glue.CfnDatabase.DatabaseInputProperty(name=GLUE_DATABASE),
)
# Phase 3 — apm_wo_snapshots table over analytics/, partitioned by dt
# with partition projection: Athena derives dt from the path, so there
# is no crawler, no MSCK REPAIR, and the classifier needs no Glue access.
analytics_location = f"s3://{self.exports_bucket.bucket_name}/analytics/"
glue.CfnTable(
self,
"SnapshotsTable",
catalog_id=self.account,
database_name=GLUE_DATABASE,
table_input=glue.CfnTable.TableInputProperty(
name=GLUE_TABLE,
table_type="EXTERNAL_TABLE",
partition_keys=[
glue.CfnTable.ColumnProperty(name="dt", type="string"),
],
parameters={
"classification": "parquet",
"EXTERNAL": "TRUE",
"projection.enabled": "true",
"projection.dt.type": "date",
"projection.dt.format": "yyyy-MM-dd",
"projection.dt.range": "2026-01-01,NOW",
"storage.location.template": f"{analytics_location}dt=${{dt}}/",
},
storage_descriptor=glue.CfnTable.StorageDescriptorProperty(
location=analytics_location,
input_format=_PARQUET_INPUT,
output_format=_PARQUET_OUTPUT,
serde_info=glue.CfnTable.SerdeInfoProperty(
serialization_library=_PARQUET_SERDE
),
columns=[
glue.CfnTable.ColumnProperty(name=name, type=type_)
for name, type_ in SNAPSHOT_COLUMNS
],
),
),
)
# Phase 3 — dedicated Athena workgroup, enforced result location + SSE-S3.
athena.CfnWorkGroup(
self,
"Workgroup",
name=ATHENA_WORKGROUP,
recursive_delete_option=True,
work_group_configuration=athena.CfnWorkGroup.WorkGroupConfigurationProperty(
enforce_work_group_configuration=True,
publish_cloud_watch_metrics_enabled=True,
result_configuration=athena.CfnWorkGroup.ResultConfigurationProperty(
output_location=f"s3://{self.exports_bucket.bucket_name}/athena-results/",
encryption_configuration=athena.CfnWorkGroup.EncryptionConfigurationProperty(
encryption_option="SSE_S3",
),
),
),
)
# Phase 2 — classifier Lambda, S3-triggered on the raw/ prefix.
# Deps (awswrangler/pandas/pyarrow/openpyxl) are Docker-bundled for ARM64
# from lambdas/classifier/requirements.txt; boto3 ships in the runtime.
@ -138,31 +231,12 @@ class PipelineStack(Stack):
s3.NotificationKeyFilter(prefix="raw/", suffix=suffix),
)
# IAM — least privilege: read raw/, read+write analytics/, catalog the
# snapshots table, and read the Anthropic key for the Haiku fallback.
# IAM — least privilege: read raw/, read+write analytics/, and read the
# Anthropic key for the Haiku fallback. No Glue access: the table is
# CDK-defined with partition projection, so the classifier only writes
# Parquet to S3 — it never touches the catalog (Phase 3).
self.exports_bucket.grant_read(self.classifier_fn, "raw/*")
self.exports_bucket.grant_read_write(self.classifier_fn, "analytics/*")
self.classifier_fn.add_to_role_policy(
iam.PolicyStatement(
sid="GlueCatalogSnapshots",
actions=[
"glue:GetDatabase",
"glue:GetTable",
"glue:CreateTable",
"glue:UpdateTable",
"glue:GetPartition",
"glue:GetPartitions",
"glue:CreatePartition",
"glue:BatchCreatePartition",
"glue:UpdatePartition",
],
resources=[
f"arn:aws:glue:{self.region}:{self.account}:catalog",
f"arn:aws:glue:{self.region}:{self.account}:database/{GLUE_DATABASE}",
f"arn:aws:glue:{self.region}:{self.account}:table/{GLUE_DATABASE}/*",
],
)
)
secretsmanager.Secret.from_secret_name_v2(
self, "AnthropicKey", ANTHROPIC_SECRET
).grant_read(self.classifier_fn)

View file

@ -4,8 +4,9 @@ Triggered on ``s3:ObjectCreated`` under the ``raw/`` prefix regardless of how th
file arrives (direct upload or local drop-folder). For each non-blank-comment row
it strips HTML, runs the two-axis classifier, derives ``is_escalation`` /
``is_action`` / ``mismatch``, writes a per-WO snapshot to ``analytics/dt=YYYY-MM-DD/``
as Parquet (registering the Glue partition), and emits a small ``summary.json``
for the slack-post Lambda to read cheaply.
as Parquet, and emits a small ``summary.json`` for the slack-post Lambda to read
cheaply. The ``apm_wo_snapshots`` Glue table is CDK-defined with partition
projection, so the write needs no Glue catalog access.
The slack-post invocation is wired in Phase 4 (the function does not exist yet).
Runtime: Python 3.12, ARM64, 512 MB, 120 s. pandas/pyarrow/awswrangler are
@ -29,8 +30,6 @@ import pandas as pd
import classify as clf
GLUE_DATABASE = "apm_wo_analysis"
GLUE_TABLE = "apm_wo_snapshots"
ANALYTICS_PREFIX = "analytics"
# Map snapshot field -> substring matched (case-insensitively) against the export
@ -202,14 +201,16 @@ def handler(event, context):
return {"classified": 0, "blank": blank}
df["dt"] = dt
# Pure Parquet write — no Glue registration. The apm_wo_snapshots table is
# CDK-defined with partition projection (Phase 3), so Athena derives the dt
# partition from the path and the classifier needs no Glue catalog access.
# overwrite_partitions keeps a same-day re-upload idempotent (replaces dt=).
wr.s3.to_parquet(
df=df,
path=f"s3://{bucket}/{ANALYTICS_PREFIX}/",
dataset=True,
partition_cols=["dt"],
mode="overwrite_partitions",
database=GLUE_DATABASE,
table=GLUE_TABLE,
)
summary = _build_summary(df, dt, key, blank)

View file

@ -0,0 +1,119 @@
"""Synth-level assertions for the Phase 3 analytics dataset.
Synthesizes ``apm-wo-analysis-pipeline`` and asserts the Glue table carries
partition projection, the column schema matches what the classifier writes, the
Athena workgroup enforces its result location, and the classifier role has zero
Glue access (projection means it never touches the catalog). No AWS, no Docker:
``aws:cdk:bundling-stacks=[]`` skips asset bundling so this is a fast offline gate.
Run with the repo venv:
python -m pytest tests/test_pipeline_synth.py -q
"""
import sys
from pathlib import Path
import aws_cdk as cdk
from aws_cdk.assertions import Match, Template
CDK_DIR = Path(__file__).resolve().parents[1] / "cdk"
sys.path.insert(0, str(CDK_DIR))
from stacks.pipeline_stack import PipelineStack # noqa: E402
def _s3_path_ending(suffix: str):
"""Match an ``Fn::Join`` S3 path (bucket name is a Ref) ending in ``suffix``."""
return {"Fn::Join": ["", Match.array_with([suffix])]}
def _template() -> Template:
app = cdk.App(context={"aws:cdk:bundling-stacks": []})
stack = PipelineStack(
app,
"apm-wo-analysis-pipeline",
env=cdk.Environment(account="328440206208", region="us-east-1"),
)
return Template.from_stack(stack)
def test_snapshots_table_uses_partition_projection():
_template().has_resource_properties(
"AWS::Glue::Table",
{
"TableInput": {
"Name": "apm_wo_snapshots",
"PartitionKeys": [{"Name": "dt", "Type": "string"}],
"Parameters": Match.object_like(
{
"projection.enabled": "true",
"projection.dt.type": "date",
"projection.dt.format": "yyyy-MM-dd",
"projection.dt.range": "2026-01-01,NOW",
"storage.location.template": _s3_path_ending(
"/analytics/dt=${dt}/"
),
}
),
}
},
)
def test_snapshots_table_column_schema_matches_classifier():
# Booleans typed correctly and the columns the BUILD.md prose omitted
# (contractor_description) are present — schema mirrors the writer.
_template().has_resource_properties(
"AWS::Glue::Table",
{
"TableInput": {
"StorageDescriptor": Match.object_like(
{
# array_with is order-sensitive: list patterns in the
# same order the classifier writes them.
"Columns": Match.array_with(
[
{"Name": "contractor_description", "Type": "string"},
{"Name": "is_escalation", "Type": "boolean"},
{"Name": "is_action", "Type": "boolean"},
{"Name": "mismatch", "Type": "string"},
]
),
}
),
}
},
)
def test_athena_workgroup_enforces_results_location():
_template().has_resource_properties(
"AWS::Athena::WorkGroup",
{
"Name": "apm-wo-analysis",
"WorkGroupConfiguration": Match.object_like(
{
"EnforceWorkGroupConfiguration": True,
"ResultConfiguration": {
"OutputLocation": _s3_path_ending("/athena-results/"),
"EncryptionConfiguration": {"EncryptionOption": "SSE_S3"},
},
}
),
},
)
def test_classifier_role_has_no_glue_access():
# Partition projection => the classifier only writes Parquet to S3. No IAM
# policy in the stack should grant any glue:* action.
policies = _template().find_resources("AWS::IAM::Policy")
for policy in policies.values():
for stmt in policy["Properties"]["PolicyDocument"]["Statement"]:
actions = stmt.get("Action", [])
actions = [actions] if isinstance(actions, str) else actions
offending = [
a for a in actions if isinstance(a, str) and a.startswith("glue:")
]
assert not offending, f"unexpected Glue access: {offending}"