Skip to content
All projects

CLOUD / DATA ENGINEERING

Care Plus AWS Data Pipeline

A serverless AWS pipeline that lands support tickets and logs in S3, cleans them with S3-triggered Lambda functions into Parquet, loads them incrementally into Redshift Serverless, and serves ad-hoc SQL in Athena and a Power BI dashboard.

  • Event-driven Lambda ETL
  • Incremental loads
  • Parquet
Context
Data engineering project
Period
2025

Stack

  • Python
  • boto3
  • pandas
  • PyArrow
  • Amazon S3
  • AWS Lambda
  • Amazon Redshift Serverless
  • Amazon Athena
  • SQL
  • Power BI

Project overview

Care Plus is a support organisation with two data sources that don't meet: support tickets in a MySQL database (priority, channel, status, agent, resolution time) and application logs written to a file per day. The pipeline brings both into one analytical store so questions like 'which channel produces the most escalations' take one SQL query.

Architecture

Raw and processed data live in separate S3 prefixes, and processing is triggered by new objects.

Implementation

  • Ingestion reads new tickets from MySQL with SQLAlchemy and uploads day-partitioned files to s3://…/raw/. A date tracker records the last loaded day, so reruns never duplicate data.
  • An S3 ObjectCreated event triggers a Lambda function per source. The log parser uses a regular expression with named groups to pull out timestamp, log level, component, event type, error flag, response time, CPU usage and user agent. The ticket cleaner trims and lower-cases fields and parses timestamps.
  • Both functions write columnar Parquet with PyArrow to a processed/ prefix, which keeps scans cheap for Athena and loads fast into Redshift.
  • Redshift Serverless loads each new day with COPY using an IAM role. Athena queries the same Parquet files directly for ad-hoc work.
Lambda entry point (simplified from the repository)python
def lambda_handler(event, context):
    record = event["Records"][0]["s3"]
    bucket, key = record["bucket"]["name"], record["object"]["key"]

    raw = read_log_from_s3(bucket, key)
    df = parse_logs(raw)                     # regex → typed columns
    out = key.replace("raw/", "processed/").replace(".log", ".parquet")
    save_parquet_to_s3(df, bucket, out)      # PyArrow → S3

Technology decisions

Event-driven Lambda instead of a scheduled cluster
The volume is a few files a day, so pay-per-invocation compute that runs only when data arrives fits better than an always-on Spark or Glue job.
Parquet in the processed layer
Columnar and compressed: Athena scans less data per query, and Redshift COPY loads it natively.
Both Athena and Redshift
Athena for exploration directly on S3, Redshift for the modelled tables behind the dashboard.

Results

The analytics layer answers ticket load by channel, status breakdown (resolved, open, escalated), daily ticket trends, error-event counts in the logs, and average CPU usage per user agent. These feed the Careplus Insights dashboard in Power BI.

What I learned

  • Keeping raw and processed data in separate prefixes makes it safe to reprocess everything after a parser change.
  • Idempotent incremental loading (tracked dates, append-only days) matters more than raw speed in small pipelines.
  • Next step: move credentials to Secrets Manager and define the bucket, functions and triggers in Terraform instead of the console.

Source