115 lines
3.5 KiB
Python
115 lines
3.5 KiB
Python
|
|
from airflow import DAG
|
||
|
|
from airflow.operators.python import PythonOperator
|
||
|
|
from datetime import datetime, timedelta
|
||
|
|
import os
|
||
|
|
|
||
|
|
# Default arguments for the DAG
|
||
|
|
default_args = {
|
||
|
|
'owner': 'airflow',
|
||
|
|
'depends_on_past': False,
|
||
|
|
'start_date': datetime(2025, 1, 1),
|
||
|
|
'email_on_failure': False,
|
||
|
|
'email_on_retry': False,
|
||
|
|
'retries': 1,
|
||
|
|
'retry_delay': timedelta(minutes=5),
|
||
|
|
}
|
||
|
|
|
||
|
|
# Define the DAG
|
||
|
|
dag = DAG(
|
||
|
|
'generate_data_all_databases',
|
||
|
|
default_args=default_args,
|
||
|
|
description='Generate fake data for all databases',
|
||
|
|
schedule=None,
|
||
|
|
catchup=False,
|
||
|
|
tags=['data', 'generation', 'fake-data'],
|
||
|
|
)
|
||
|
|
|
||
|
|
# Script paths
|
||
|
|
SCRIPTS_DIR = '/root/airflow/dags/scripts'
|
||
|
|
POSTGRES_SCRIPT = os.path.join(SCRIPTS_DIR, 'generate_postgres_sales_data.py')
|
||
|
|
MONGODB_SCRIPT = os.path.join(SCRIPTS_DIR, 'generate_mongodb_events_data.py')
|
||
|
|
MYSQL_SCRIPT = os.path.join(SCRIPTS_DIR, 'generate_mysql_employee_data.py')
|
||
|
|
NEO4J_SCRIPT = os.path.join(SCRIPTS_DIR, 'generate_neo4j_graph_data.py')
|
||
|
|
CASSANDRA_SCRIPT = os.path.join(SCRIPTS_DIR, 'generate_cassandra_telemetry_data.py')
|
||
|
|
|
||
|
|
# Task functions
|
||
|
|
def run_postgres_script():
|
||
|
|
import subprocess
|
||
|
|
result = subprocess.run(['python3', POSTGRES_SCRIPT], capture_output=True, text=True)
|
||
|
|
print(result.stdout)
|
||
|
|
if result.stderr:
|
||
|
|
print(result.stderr)
|
||
|
|
if result.returncode != 0:
|
||
|
|
raise Exception(f"PostgreSQL script failed with return code {result.returncode}")
|
||
|
|
|
||
|
|
def run_mongodb_script():
|
||
|
|
import subprocess
|
||
|
|
result = subprocess.run(['python3', MONGODB_SCRIPT], capture_output=True, text=True)
|
||
|
|
print(result.stdout)
|
||
|
|
if result.stderr:
|
||
|
|
print(result.stderr)
|
||
|
|
if result.returncode != 0:
|
||
|
|
raise Exception(f"MongoDB script failed with return code {result.returncode}")
|
||
|
|
|
||
|
|
def run_mysql_script():
|
||
|
|
import subprocess
|
||
|
|
result = subprocess.run(['python3', MYSQL_SCRIPT], capture_output=True, text=True)
|
||
|
|
print(result.stdout)
|
||
|
|
if result.stderr:
|
||
|
|
print(result.stderr)
|
||
|
|
if result.returncode != 0:
|
||
|
|
raise Exception(f"MySQL script failed with return code {result.returncode}")
|
||
|
|
|
||
|
|
def run_neo4j_script():
|
||
|
|
import subprocess
|
||
|
|
result = subprocess.run(['python3', NEO4J_SCRIPT], capture_output=True, text=True)
|
||
|
|
print(result.stdout)
|
||
|
|
if result.stderr:
|
||
|
|
print(result.stderr)
|
||
|
|
if result.returncode != 0:
|
||
|
|
raise Exception(f"Neo4j script failed with return code {result.returncode}")
|
||
|
|
|
||
|
|
def run_cassandra_script():
|
||
|
|
import subprocess
|
||
|
|
result = subprocess.run(['python3', CASSANDRA_SCRIPT], capture_output=True, text=True)
|
||
|
|
print(result.stdout)
|
||
|
|
if result.stderr:
|
||
|
|
print(result.stderr)
|
||
|
|
if result.returncode != 0:
|
||
|
|
raise Exception(f"Cassandra script failed with return code {result.returncode}")
|
||
|
|
|
||
|
|
# Define tasks
|
||
|
|
generate_postgres_task = PythonOperator(
|
||
|
|
task_id='generate_postgres_sales_data',
|
||
|
|
python_callable=run_postgres_script,
|
||
|
|
dag=dag,
|
||
|
|
)
|
||
|
|
|
||
|
|
generate_mongodb_task = PythonOperator(
|
||
|
|
task_id='generate_mongodb_events_data',
|
||
|
|
python_callable=run_mongodb_script,
|
||
|
|
dag=dag,
|
||
|
|
)
|
||
|
|
|
||
|
|
generate_mysql_task = PythonOperator(
|
||
|
|
task_id='generate_mysql_employee_data',
|
||
|
|
python_callable=run_mysql_script,
|
||
|
|
dag=dag,
|
||
|
|
)
|
||
|
|
|
||
|
|
generate_neo4j_task = PythonOperator(
|
||
|
|
task_id='generate_neo4j_graph_data',
|
||
|
|
python_callable=run_neo4j_script,
|
||
|
|
dag=dag,
|
||
|
|
)
|
||
|
|
|
||
|
|
generate_cassandra_task = PythonOperator(
|
||
|
|
task_id='generate_cassandra_telemetry_data',
|
||
|
|
python_callable=run_cassandra_script,
|
||
|
|
dag=dag,
|
||
|
|
)
|
||
|
|
|
||
|
|
# Set task dependencies - run all in parallel
|
||
|
|
[generate_postgres_task, generate_mongodb_task, generate_mysql_task,
|
||
|
|
generate_neo4j_task, generate_cassandra_task]
|