Merge pull request #9 from Sea-Haven-Industries/feature/phase-3-analytics
Some checks are pending
Deploy / deploy (push) Waiting to run

Phase 3: analytics dataset (Glue projection table + Athena)
This commit is contained in:
Adam Moussa 2026-05-29 13:49:47 -04:00 • committed by GitHub
commit 4cc58f5683
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 235 additions and 38 deletions

View file

@ -142,7 +142,7 @@ cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role
## Status
**Phase 2 — classifier (in review).** Build-out proceeds per
**Phase 3 — analytics dataset (in review).** Build-out proceeds per
[`docs/BUILD.md`](./docs/BUILD.md): ingestion → classifier → Glue/Athena → Slack
→ Grafana → docs.
@ -151,5 +151,8 @@ cdk deploy apm-wo-analysis-grafana # EC2, ALB, SG, Route53, datasource role
deployed; PR #7 open.
- **Phase 2** classifier Lambda + Glue database + S3 `raw/` trigger — implemented
and `cdk synth`-green; PR #8 open (stacked on Phase 1, not yet deployed).
- **Phases 3–6** (Glue/Athena query layer, Slack, Grafana, final docs) — not
started.
- **Phase 3** `apm_wo_snapshots` projection table + Athena workgroup —
implemented and `cdk synth`-green; PR #9 open (stacked on Phase 2). The
classifier writes pure Parquet and holds no Glue access (projection handles
partitions).
- **Phases 4–6** (Slack, Grafana, final docs) — not started.

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}"