airflow-dag
Create Apache Airflow DAGs for construction data pipelines. Orchestrate ETL, validation, and reporting workflows.
What this skill does
# Apache Airflow DAG for Construction
## Overview
Apache Airflow orchestrates complex data pipelines. This skill creates DAGs for construction ETL processes - from BIM extraction to cost reports.
## Python Implementation
```python
from datetime import datetime, timedelta
from typing import Dict, Any, List, Optional, Callable
from dataclasses import dataclass
from enum import Enum
import json
class TaskStatus(Enum):
"""Task execution status."""
PENDING = "pending"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
SKIPPED = "skipped"
@dataclass
class DAGTask:
"""Single task in DAG."""
task_id: str
operator: str
params: Dict[str, Any]
upstream: List[str]
downstream: List[str]
@dataclass
class DAGConfig:
"""DAG configuration."""
dag_id: str
schedule: str
start_date: datetime
catchup: bool
default_args: Dict[str, Any]
tags: List[str]
class ConstructionDAGBuilder:
"""Build Airflow DAGs for construction pipelines."""
# Default DAG arguments
DEFAULT_ARGS = {
'owner': 'ddc',
'depends_on_past': False,
'email_on_failure': True,
'email_on_retry': False,
'retries': 2,
'retry_delay': timedelta(minutes=5),
'execution_timeout': timedelta(hours=2)
}
def __init__(self, dag_id: str,
schedule: str = '@daily',
tags: List[str] = None):
self.dag_id = dag_id
self.schedule = schedule
self.tags = tags or ['construction', 'ddc']
self.tasks: Dict[str, DAGTask] = {}
def add_bash_task(self, task_id: str,
command: str,
upstream: List[str] = None) -> str:
"""Add bash command task."""
self.tasks[task_id] = DAGTask(
task_id=task_id,
operator='BashOperator',
params={'bash_command': command},
upstream=upstream or [],
downstream=[]
)
self._update_downstream(task_id, upstream)
return task_id
def add_python_task(self, task_id: str,
python_callable: str,
op_kwargs: Dict = None,
upstream: List[str] = None) -> str:
"""Add Python callable task."""
self.tasks[task_id] = DAGTask(
task_id=task_id,
operator='PythonOperator',
params={
'python_callable': python_callable,
'op_kwargs': op_kwargs or {}
},
upstream=upstream or [],
downstream=[]
)
self._update_downstream(task_id, upstream)
return task_id
def add_sensor_task(self, task_id: str,
filepath: str,
upstream: List[str] = None) -> str:
"""Add file sensor task."""
self.tasks[task_id] = DAGTask(
task_id=task_id,
operator='FileSensor',
params={
'filepath': filepath,
'poke_interval': 300,
'timeout': 3600
},
upstream=upstream or [],
downstream=[]
)
self._update_downstream(task_id, upstream)
return task_id
def add_branch_task(self, task_id: str,
python_callable: str,
upstream: List[str] = None) -> str:
"""Add branching task."""
self.tasks[task_id] = DAGTask(
task_id=task_id,
operator='BranchPythonOperator',
params={'python_callable': python_callable},
upstream=upstream or [],
downstream=[]
)
self._update_downstream(task_id, upstream)
return task_id
def _update_downstream(self, task_id: str, upstream: List[str]):
"""Update downstream references."""
if upstream:
for up_task in upstream:
if up_task in self.tasks:
self.tasks[up_task].downstream.append(task_id)
def generate_dag_code(self) -> str:
"""Generate Airflow DAG Python code."""
code = '''
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.sensors.filesystem import FileSensor
from datetime import datetime, timedelta
default_args = {
'owner': 'ddc',
'depends_on_past': False,
'email_on_failure': True,
'retries': 2,
'retry_delay': timedelta(minutes=5),
}
'''
code += f'''
with DAG(
dag_id='{self.dag_id}',
default_args=default_args,
schedule_interval='{self.schedule}',
start_date=datetime(2024, 1, 1),
catchup=False,
tags={self.tags}
) as dag:
'''
# Generate task definitions
for task_id, task in self.tasks.items():
code += self._generate_task_code(task)
code += '\n'
# Generate dependencies
code += '\n # Task dependencies\n'
for task_id, task in self.tasks.items():
if task.upstream:
for upstream in task.upstream:
code += f" {upstream} >> {task_id}\n"
return code
def _generate_task_code(self, task: DAGTask) -> str:
"""Generate code for single task."""
if task.operator == 'BashOperator':
return f''' {task.task_id} = BashOperator(
task_id='{task.task_id}',
bash_command="{task.params['bash_command']}"
)'''
elif task.operator == 'PythonOperator':
kwargs = json.dumps(task.params.get('op_kwargs', {}))
return f''' {task.task_id} = PythonOperator(
task_id='{task.task_id}',
python_callable={task.params['python_callable']},
op_kwargs={kwargs}
)'''
elif task.operator == 'FileSensor':
return f''' {task.task_id} = FileSensor(
task_id='{task.task_id}',
filepath='{task.params["filepath"]}',
poke_interval={task.params['poke_interval']},
timeout={task.params['timeout']}
)'''
elif task.operator == 'BranchPythonOperator':
return f''' {task.task_id} = BranchPythonOperator(
task_id='{task.task_id}',
python_callable={task.params['python_callable']}
)'''
return ''
def save_dag(self, output_path: str):
"""Save DAG to file."""
code = self.generate_dag_code()
with open(output_path, 'w') as f:
f.write(code)
return output_path
class ConstructionPipelineTemplates:
"""Pre-built construction pipeline templates."""
@staticmethod
def bim_validation_pipeline(dag_id: str = 'bim_validation') -> ConstructionDAGBuilder:
"""Create BIM validation pipeline."""
builder = ConstructionDAGBuilder(dag_id, schedule='@daily',
tags=['bim', 'validation'])
# Wait for file
builder.add_sensor_task('wait_for_model', '/data/input/*.ifc')
# Convert to Excel
builder.add_bash_task(
'convert_ifc',
'IfcExporter.exe /data/input/*.ifc bbox',
upstream=['wait_for_model']
)
# Validate data
builder.add_python_task(
'validate_data',
'validate_bim_data',
{'rules_file': '/config/validation_rules.xlsx'},
upstream=['convert_ifc']
)
# Branch based on validation
builder.add_branch_task(
'check_validation',
'check_validation_result',
upstream=['validate_data']
)
# Success path
builder.add_python_task(
'generate_report',
'generate_validation_report',
upstream=['check_validation']
)
# Failure path
builder.add_python_task(
'send_alert',
'send_validation_alert',
upstream=['check_valiRelated in Data & Analytics
clawarr-suite
IncludedComprehensive management for self-hosted media stacks (Sonarr, Radarr, Lidarr, Readarr, Prowlarr, Bazarr, Overseerr, Plex, Tautulli, SABnzbd, Recyclarr, Unpackerr, Notifiarr, Maintainerr, Kometa, FlareSolverr). Deep library exploration, analytics, dashboard generation, content management, request handling, subtitle management, indexer control, download monitoring, quality profile sync, library cleanup automation, notification routing, collection/overlay management, and media tracker integration (Trakt, Letterboxd, Simkl).
querying-soql
IncludedSOQL query generation, optimization, and analysis with 100-point scoring. Use this skill when the user needs SOQL/SOSL authoring or optimization: natural-language-to-query generation, relationship queries, aggregates, query-plan analysis, and performance or safety improvements for Salesforce queries. TRIGGER when: user writes, optimizes, or debugs SOQL/SOSL queries, touches .soql files, or asks about relationship queries, aggregates, or query performance. DO NOT TRIGGER when: bulk data operations (use handling-sf-data), Apex DML logic (use generating-apex), or report/dashboard queries.
app-store-optimization
IncludedApp Store Optimization (ASO) toolkit for researching keywords, analyzing competitor rankings, generating metadata suggestions, and improving app visibility on Apple App Store and Google Play Store. Use when the user asks about ASO, app store rankings, app metadata, app titles and descriptions, app store listings, app visibility, or mobile app marketing on iOS or Android. Supports keyword research and scoring, competitor keyword analysis, metadata optimization, A/B test planning, launch checklists, and tracking ranking changes.
habit-flow
IncludedAI-powered atomic habit tracker with natural language logging, streak tracking, smart reminders, and coaching. Use for creating habits, logging completions naturally ("I meditated today"), viewing progress, and getting personalized coaching.
app-store-optimization
IncludedApp Store Optimization (ASO) toolkit for researching keywords, analyzing competitor rankings, generating metadata suggestions, and improving app visibility on Apple App Store and Google Play Store. Use when the user asks about ASO, app store rankings, app metadata, app titles and descriptions, app store listings, app visibility, or mobile app marketing on iOS or Android. Supports keyword research and scoring, competitor keyword analysis, metadata optimization, A/B test planning, launch checklists, and tracking ranking changes.
visualizing-data
IncludedBuilds dashboards, reports, and data-driven interfaces requiring charts, graphs, or visual analytics. Provides systematic framework for selecting appropriate visualizations based on data characteristics and analytical purpose. Includes 24+ visualization types organized by purpose (trends, comparisons, distributions, relationships, flows, hierarchies, geospatial), accessibility patterns (WCAG 2.1 AA compliance), colorblind-safe palettes, and performance optimization strategies. Use when creating visualizations, choosing chart types, displaying data graphically, or designing data interfaces.