Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions examples/bitcoin-tracker/AIRFLOW_INTEGRATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,10 @@ sys.path.insert(0, os.path.abspath(os.path.dirname(__file__) + "/../.."))

from ingest_bitcoin_prices import fetch_bitcoin_price, insert_to_bigquery

# Configuration from environment variables
GCP_PROJECT_ID = os.environ["GCP_PROJECT_ID"] # required: fail at parse time rather than
# targeting a placeholder project that does not exist

# Default arguments
default_args = {
"owner": "data-engineering",
Expand Down Expand Up @@ -256,7 +260,7 @@ with DAG(
"""Insert price data to BigQuery."""
price_data = context["ti"].xcom_pull(task_ids="fetch_bitcoin_price")

project_id = os.environ.get("GCP_PROJECT_ID", "<<YOUR_PROJECT_HERE>>")
project_id = GCP_PROJECT_ID
dataset = "crypto_data"
table = "bitcoin_prices"

Expand Down Expand Up @@ -286,7 +290,7 @@ with DAG(
task_id="check_data_quality",
sql=f"""
SELECT COUNT(*) > 0
FROM `{os.environ.get('GCP_PROJECT_ID', '<<YOUR_PROJECT_HERE>>')}.crypto_data.bitcoin_prices`
FROM `{GCP_PROJECT_ID}.crypto_data.bitcoin_prices`
WHERE DATE(timestamp) = CURRENT_DATE()
""",
use_legacy_sql=False,
Expand All @@ -297,7 +301,7 @@ with DAG(
task_id="verify_transformations",
sql=f"""
SELECT COUNT(*) > 0
FROM `{os.environ.get('GCP_PROJECT_ID', '<<YOUR_PROJECT_HERE>>')}.crypto_data.daily_price_summary`
FROM `{GCP_PROJECT_ID}.crypto_data.daily_price_summary`
WHERE price_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 1 DAY)
""",
use_legacy_sql=False,
Expand Down
7 changes: 6 additions & 1 deletion examples/bitcoin-tracker/airflow-quickstart.sh
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,12 @@ NC='\033[0m' # No Color
# Configuration
PROJECT_DIR=$(pwd)
AIRFLOW_HOME=${AIRFLOW_HOME:-$HOME/airflow}
GCP_PROJECT_ID=${GCP_PROJECT_ID:-<<YOUR_PROJECT_HERE>>}
GCP_PROJECT_ID=${GCP_PROJECT_ID:-}
if [ -z "$GCP_PROJECT_ID" ]; then
echo "❌ Error: GCP_PROJECT_ID environment variable not set"
echo "Usage: export GCP_PROJECT_ID=your-project-id && ./airflow-quickstart.sh"
exit 1
fi

echo -e "${BLUE}Configuration:${NC}"
echo " Project Directory: $PROJECT_DIR"
Expand Down
13 changes: 12 additions & 1 deletion examples/bitcoin-tracker/load_bitcoin_price_batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from google.cloud import bigquery
from datetime import datetime
import os
import sys
import json
import tempfile

Expand Down Expand Up @@ -67,7 +68,17 @@ def load_to_bigquery_batch(row, project_id, dataset_id="crypto_data", table_id="
os.unlink(temp_file)

if __name__ == "__main__":
project_id = os.getenv("GCP_PROJECT_ID", "<<YOUR_PROJECT_HERE>>")
# Get project ID from environment or command line
project_id = os.getenv("GCP_PROJECT_ID")

if not project_id and len(sys.argv) > 1:
project_id = sys.argv[1]

if not project_id:
print("❌ Error: GCP_PROJECT_ID environment variable not set")
print("Usage: python load_bitcoin_price_batch.py [PROJECT_ID]")
print(" or: export GCP_PROJECT_ID=your-project-id && python load_bitcoin_price_batch.py")
sys.exit(1)

print(f"🚀 Fetching Bitcoin price...")
price_data = fetch_bitcoin_price()
Expand Down
14 changes: 13 additions & 1 deletion examples/bitcoin-tracker/runtime/ingest_bitcoin_prices.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,19 @@

def main():
"""Fetch Bitcoin price from CoinGecko API and insert to BigQuery"""
# Resolve the target project before doing any work, so a missing
# GCP_PROJECT_ID fails here rather than silently targeting a placeholder.
project_id = os.getenv("GCP_PROJECT_ID")

if not project_id and len(sys.argv) > 1:
project_id = sys.argv[1]

if not project_id:
print("❌ Error: GCP_PROJECT_ID environment variable not set")
print("Usage: python ingest_bitcoin_prices.py [PROJECT_ID]")
print(" or: export GCP_PROJECT_ID=your-project-id && python ingest_bitcoin_prices.py")
sys.exit(1)

print("🚀 Starting Bitcoin price ingestion...")

# Fetch from CoinGecko API (free tier, no auth required)
Expand Down Expand Up @@ -47,7 +60,6 @@ def main():
}

# Insert to BigQuery
project_id = os.getenv("GCP_PROJECT_ID", "<<YOUR_PROJECT_HERE>>")
dataset_id = "crypto_data"
table_id = "bitcoin_prices"

Expand Down