| name | harvard-artifacts-collection-analytics-pipeline |
| description | End-to-end data engineering pipeline for Harvard Art Museums API with ETL, SQL analytics, and Streamlit visualization |
| triggers | ["build a data pipeline for Harvard Art Museums API","create ETL workflow for museum artifacts data","set up artifact collection analytics with Streamlit","implement Harvard museums data engineering pipeline","analyze art museum data with SQL and visualization","extract and transform Harvard API artifact data","build interactive dashboard for museum collection data","create SQL analytics for art artifacts"] |
Harvard Artifacts Collection Analytics Pipeline
Skill by ara.so — Data Skills collection.
Overview
This project provides a complete data engineering solution for the Harvard Art Museums API, featuring:
- ETL pipeline for artifact metadata, media, and color data
- SQL database storage (MySQL/TiDB Cloud)
- 20+ analytical SQL queries
- Interactive Streamlit dashboard with Plotly visualizations
The architecture follows: API → ETL → SQL → Analytics → Visualization
Installation
git clone https://github.com/Manali0711/Harvard-Artifacts-Collection-Data-Engineering-Analytics-App.git
cd Harvard-Artifacts-Collection-Data-Engineering-Analytics-App
pip install -r requirements.txt
Required Dependencies
streamlit
pandas
requests
mysql-connector-python
plotly
python-dotenv
Configuration
Environment Variables
Create a .env file in the project root:
HARVARD_API_KEY=your_api_key_here
DB_HOST=your_database_host
DB_PORT=3306
DB_USER=your_username
DB_PASSWORD=your_password
DB_NAME=harvard_artifacts
Database Setup
import mysql.connector
from mysql.connector import Error
def create_database_connection():
"""Establish MySQL/TiDB connection"""
try:
connection = mysql.connector.connect(
host=os.getenv('DB_HOST'),
port=os.getenv('DB_PORT'),
user=os.getenv('DB_USER'),
password=os.getenv('DB_PASSWORD'),
database=os.getenv('DB_NAME')
)
return connection
except Error as e:
print(f"Database connection error: {e}")
return None
def create_tables(connection):
"""Create database schema"""
cursor = connection.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS artifactmetadata (
artifact_id INT PRIMARY KEY,
title VARCHAR(500),
culture VARCHAR(200),
century VARCHAR(100),
classification VARCHAR(200),
department VARCHAR(200),
dated VARCHAR(200),
period VARCHAR(200),
technique VARCHAR(500),
medium VARCHAR(500),
dimensions VARCHAR(500),
creditline TEXT,
url VARCHAR(500)
)
""")
cursor.execute("""
CREATE TABLE IF NOT EXISTS artifactmedia (
media_id INT AUTO_INCREMENT PRIMARY KEY,
artifact_id INT,
image_url VARCHAR(1000),
caption TEXT,
FOREIGN KEY (artifact_id) REFERENCES artifactmetadata(artifact_id)
)
""")
cursor.execute("""
CREATE TABLE IF NOT EXISTS artifactcolors (
color_id INT AUTO_INCREMENT PRIMARY KEY,
artifact_id INT,
color_hex VARCHAR(10),
color_name VARCHAR(100),
percentage FLOAT,
FOREIGN KEY (artifact_id) REFERENCES artifactmetadata(artifact_id)
)
""")
connection.commit()
cursor.close()
ETL Pipeline
Extract: Fetch Data from Harvard API
import requests
import time
def fetch_artifacts_from_api(api_key, size=100, page=1):
"""Extract artifacts from Harvard Art Museums API"""
base_url = "https://api.harvardartmuseums.org/object"
params = {
'apikey': api_key,
'size': size,
'page': page,
'hasimage': 1
}
try:
response = requests.get(base_url, params=params)
response.raise_for_status()
data = response.json()
time.sleep(0.5)
return data.get('records', []), data.get('info', {})
except requests.exceptions.RequestException as e:
print(f"API request error: {e}")
return [], {}
def paginate_api_collection(api_key, max_pages=10):
"""Collect multiple pages of artifacts"""
all_artifacts = []
for page in range(1, max_pages + 1):
records, info = fetch_artifacts_from_api(api_key, page=page)
if not records:
break
all_artifacts.extend(records)
print(f"Fetched page {page}, total artifacts: ")
all_artifacts
Transform: Process JSON Data
import pandas as pd
def transform_artifact_metadata(artifacts):
"""Transform artifact data into structured format"""
metadata = []
for artifact in artifacts:
metadata.append({
'artifact_id': artifact.get('id'),
'title': artifact.get('title'),
'culture': artifact.get('culture'),
'century': artifact.get('century'),
'classification': artifact.get('classification'),
'department': artifact.get('department'),
'dated': artifact.get('dated'),
'period': artifact.get('period'),
'technique': artifact.get('technique'),
'medium': artifact.get('medium'),
'dimensions': artifact.get('dimensions'),
'creditline': artifact.get('creditline'),
'url': artifact.get('url')
})
return pd.DataFrame(metadata)
def transform_artifact_media(artifacts):
"""Extract media/image data"""
media_data = []
for artifact in artifacts:
artifact_id = artifact.get('id')
images = artifact.get('images', [])
image images:
media_data.append({
: artifact_id,
: image.get(),
: image.get()
})
pd.DataFrame(media_data)
():
color_data = []
artifact artifacts:
artifact_id = artifact.get()
colors = artifact.get(, [])
color colors:
color_data.append({
: artifact_id,
: color.get(),
: color.get(),
: color.get()
})
pd.DataFrame(color_data)
Load: Insert into Database
def load_dataframe_to_sql(df, table_name, connection):
"""Batch insert DataFrame into SQL table"""
cursor = connection.cursor()
columns = ', '.join(df.columns)
placeholders = ', '.join(['%s'] * len(df.columns))
insert_query = f"INSERT IGNORE INTO {table_name} ({columns}) VALUES ({placeholders})"
data_tuples = [tuple(row) for row in df.values]
cursor.executemany(insert_query, data_tuples)
connection.commit()
cursor.close()
print(f"Inserted {len(df)} records into {table_name}")
def run_etl_pipeline(api_key, connection, max_pages=5):
"""Execute complete ETL pipeline"""
artifacts = paginate_api_collection(api_key, max_pages)
metadata_df = transform_artifact_metadata(artifacts)
media_df = transform_artifact_media(artifacts)
colors_df = transform_artifact_colors(artifacts)
load_dataframe_to_sql(metadata_df, 'artifactmetadata', connection)
load_dataframe_to_sql(media_df, 'artifactmedia', connection)
load_dataframe_to_sql(colors_df, 'artifactcolors', connection)
return len(artifacts)
SQL Analytics Queries
Sample Analytical Queries
ANALYTICAL_QUERIES = {
"Artifacts by Culture": """
SELECT culture, COUNT(*) as artifact_count
FROM artifactmetadata
WHERE culture IS NOT NULL
GROUP BY culture
ORDER BY artifact_count DESC
LIMIT 10
""",
"Artifacts by Century": """
SELECT century, COUNT(*) as count
FROM artifactmetadata
WHERE century IS NOT NULL
GROUP BY century
ORDER BY count DESC
""",
"Department Distribution": """
SELECT department, COUNT(*) as total_artifacts
FROM artifactmetadata
GROUP BY department
ORDER BY total_artifacts DESC
""",
"Most Common Colors": """
SELECT color_name, COUNT(*) as usage_count, AVG(percentage) as avg_percentage
FROM artifactcolors
WHERE color_name IS NOT NULL
GROUP BY color_name
ORDER BY usage_count DESC
LIMIT 15
""",
"Media Availability": """
SELECT
COUNT(DISTINCT m.artifact_id) as artifacts_with_media,
COUNT(*) as total_images
FROM artifactmedia m
""",
"Classification Analysis": """
SELECT classification, COUNT(*) as count,
GROUP_CONCAT(DISTINCT culture SEPARATOR ', ') as cultures
FROM artifactmetadata
WHERE classification IS NOT NULL
GROUP BY classification
ORDER BY count DESC
LIMIT 10
"""
}
def execute_query(connection, query_name):
"""Run analytical query and return results"""
cursor = connection.cursor(dictionary=True)
query = ANALYTICAL_QUERIES[query_name]
cursor.execute(query)
results = cursor.fetchall()
cursor.close()
return pd.DataFrame(results)
Streamlit Dashboard
Main Application Structure
import streamlit as st
import plotly.express as px
import os
from dotenv import load_dotenv
load_dotenv()
def main():
st.set_page_config(
page_title="Harvard Artifacts Analytics",
page_icon="🏛️",
layout="wide"
)
st.title("🏛️ Harvard Art Museums Analytics Dashboard")
st.markdown("---")
with st.sidebar:
st.header("⚙️ Configuration")
api_key = st.text_input(
"Harvard API Key",
value=os.getenv('HARVARD_API_KEY', ''),
type="password"
)
if st.button("Connect to Database"):
connection = create_database_connection()
if connection:
st.success("✅ Database connected!")
st.session_state['db_connection'] = connection
else:
st.error("❌ Connection failed")
tab1, tab2, tab3 = st.tabs(["📥 ETL Pipeline", "📊 Analytics", "📈 Visualizations"])
with tab1:
render_etl_tab(api_key)
with tab2:
render_analytics_tab()
with tab3:
render_visualization_tab()
def render_etl_tab():
st.header()
col1, col2 = st.columns()
col1:
max_pages = st.slider(, , , )
col2:
st.button(, =):
api_key:
st.error()
connection = st.session_state.get()
connection:
st.error()
st.spinner():
:
create_tables(connection)
total_artifacts = run_etl_pipeline(api_key, connection, max_pages)
st.success()
Exception e:
st.error()
():
st.header()
connection = st.session_state.get()
connection:
st.warning()
selected_query = st.selectbox(
,
(ANALYTICAL_QUERIES.keys())
)
st.button():
st.spinner():
:
df_results = execute_query(connection, selected_query)
st.subheader()
st.dataframe(df_results, use_container_width=)
(df_results) > :
st.session_state[] = df_results
st.session_state[] = selected_query
Exception e:
st.error()
():
st.header()
st.session_state:
st.info()
df = st.session_state[]
query_name = st.session_state[]
(df.columns) >= :
x_col = df.columns[]
y_col = df.columns[]
fig = px.bar(
df,
x=x_col,
y=y_col,
title=query_name,
template=
)
st.plotly_chart(fig, use_container_width=)
__name__ == :
main()
Running the Dashboard
streamlit run app.py
Common Patterns
Incremental Data Loading
def get_max_artifact_id(connection):
"""Get highest artifact ID in database"""
cursor = connection.cursor()
cursor.execute("SELECT MAX(artifact_id) FROM artifactmetadata")
result = cursor.fetchone()
cursor.close()
return result[0] or 0
def incremental_etl(api_key, connection):
"""Load only new artifacts"""
max_id = get_max_artifact_id(connection)
params = {'apikey': api_key, 'q': f'id:>{max_id}'}
Error Handling and Logging
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
def safe_etl_execution(api_key, connection, max_pages):
"""ETL with comprehensive error handling"""
try:
artifacts = paginate_api_collection(api_key, max_pages)
logger.info(f"Extracted {len(artifacts)} artifacts")
metadata_df = transform_artifact_metadata(artifacts)
logger.info(f"Transformed {len(metadata_df)} metadata records")
load_dataframe_to_sql(metadata_df, 'artifactmetadata', connection)
logger.info("Successfully loaded to database")
return True
except Exception as e:
logger.error(f"ETL pipeline failed: {e}")
return False
Troubleshooting
API Rate Limiting
import time
from functools import wraps
def retry_with_backoff(retries=3, backoff_in_seconds=1):
def decorator(func):
@wraps(func)
def wrapper(*args, **kwargs):
x = 0
while x < retries:
try:
return func(*args, **kwargs)
except requests.exceptions.HTTPError as e:
if e.response.status_code == 429:
sleep_time = backoff_in_seconds * (2 ** x)
time.sleep(sleep_time)
x += 1
else:
raise
return func(*args, **kwargs)
return wrapper
return decorator
@retry_with_backoff(retries=5)
def fetch_with_retry(url, params):
response = requests.get(url, params=params)
response.raise_for_status()
return response.json()
Database Connection Issues
from mysql.connector import pooling
def create_connection_pool():
"""Create reusable connection pool"""
return pooling.MySQLConnectionPool(
pool_name="harvard_pool",
pool_size=5,
host=os.getenv('DB_HOST'),
port=os.getenv('DB_PORT'),
user=os.getenv('DB_USER'),
password=os.getenv('DB_PASSWORD'),
database=os.getenv('DB_NAME')
)
pool = create_connection_pool()
connection = pool.get_connection()
Memory Management for Large Datasets
def chunked_data_load(artifacts, chunk_size=100):
"""Process large datasets in chunks"""
for i in range(0, len(artifacts), chunk_size):
chunk = artifacts[i:i + chunk_size]
metadata_df = transform_artifact_metadata(chunk)
load_dataframe_to_sql(metadata_df, 'artifactmetadata', connection)
del metadata_df
Key Features Summary
- ETL Pipeline: Automated data collection with pagination and rate limiting
- SQL Storage: Normalized schema with foreign key relationships
- Analytics: 20+ pre-built queries for artifact insights
- Visualization: Interactive Plotly charts in Streamlit
- Scalability: Handles batch processing and incremental loads