apm-wo-analysis/cdk/stacks/pipeline_stack.py

245 lines
9.4 KiB
Python
Raw Normal View History

"""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.
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.
No DynamoDB — this is an S3 + Athena analytics workload (see CLAUDE.md).
"""
import os
from aws_cdk import (
BundlingOptions,
Duration,
RemovalPolicy,
Stack,
)
from aws_cdk import (
aws_athena as athena,
)
from aws_cdk import (
aws_glue as glue,
)
from aws_cdk import (
aws_iam as iam,
)
from aws_cdk import (
aws_lambda as lambda_,
)
from aws_cdk import (
aws_logs as logs,
)
from aws_cdk import (
aws_s3 as s3,
)
from aws_cdk import (
aws_s3_notifications as s3n,
)
from aws_cdk import (
aws_secretsmanager as secretsmanager,
)
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:
super().__init__(scope, construct_id, **kwargs)
# Phase 1 — single exports bucket.
# Prefixes: raw/ (incoming), analytics/ (per-WO snapshots), athena-results/.
self.exports_bucket = s3.Bucket(
self,
"Exports",
bucket_name=f"apm-wo-analysis-exports-{self.account}",
encryption=s3.BucketEncryption.S3_MANAGED,
block_public_access=s3.BlockPublicAccess.BLOCK_ALL,
enforce_ssl=True,
removal_policy=RemovalPolicy.RETAIN,
lifecycle_rules=[
s3.LifecycleRule(
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),
),
],
)
# Phase 1 — least-privilege identity for the local drop-folder uploader.
# Scoped to s3:PutObject on raw/* only. The access key is created
# out-of-band (aws iam create-access-key) and stored in the local
# ~/.aws/credentials profile `apm-wo-drop` — never in CloudFormation.
self.drop_uploader = iam.User(
self, "DropUploader", user_name="apm-wo-drop-uploader"
)
self.drop_uploader.add_to_policy(
iam.PolicyStatement(
sid="PutRawExportsOnly",
actions=["s3:PutObject"],
resources=[self.exports_bucket.arn_for_objects("raw/*")],
)
)
# Phase 2/3 — Glue database for the analytics dataset.
glue.CfnDatabase(
self,
"AnalyticsDb",
catalog_id=self.account,
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.
classifier_logs = logs.LogGroup(
self,
"ClassifierLogs",
log_group_name="/aws/lambda/apm-wo-analysis-classifier",
retention=logs.RetentionDays.TWO_MONTHS,
removal_policy=RemovalPolicy.DESTROY,
)
self.classifier_fn = lambda_.Function(
self,
"Classifier",
function_name="apm-wo-analysis-classifier",
runtime=lambda_.Runtime.PYTHON_3_12,
architecture=lambda_.Architecture.ARM_64,
handler="handler.handler",
memory_size=512,
timeout=Duration.seconds(120),
log_group=classifier_logs,
environment={"APM_HAIKU_FALLBACK": "on"},
code=lambda_.Code.from_asset(
os.path.join(LAMBDAS_DIR, "classifier"),
bundling=BundlingOptions(
image=lambda_.Runtime.PYTHON_3_12.bundling_image,
platform="linux/arm64",
command=[
"bash",
"-c",
"pip install -r requirements.txt -t /asset-output "
"&& cp -au . /asset-output",
],
),
),
)
# S3 trigger: any .xlsx/.csv landing under raw/ invokes the classifier.
for suffix in (".xlsx", ".csv"):
self.exports_bucket.add_event_notification(
s3.EventType.OBJECT_CREATED,
s3n.LambdaDestination(self.classifier_fn),
s3.NotificationKeyFilter(prefix="raw/", suffix=suffix),
)
# 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/*")
secretsmanager.Secret.from_secret_name_v2(
self, "AnthropicKey", ANTHROPIC_SECRET
).grant_read(self.classifier_fn)
# Phase 4 — slack-post Lambda + scoped IAM; classifier async-invokes it. TODO