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
- Links
- care-plus-de
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
Source
- MySQLsupport tickets
- Log filesone per day
Landing
- S3 raw/CSV and .log
Transform
- AWS LambdaS3 event trigger
- pandas · regex
Processed
- S3 processed/Parquet
Analytics
- Redshift ServerlessCOPY, incremental
- Athenaad-hoc SQL
- Power BIdashboard
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.
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 → S3Technology 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.