用 Codex 或 Claude 帮你安装 复制这段 Prompt,粘贴到 Codex、Claude 或其他助手里,让它检查 Skill 页面并帮你完成安装。
直接命令不会经过审查 Prompt;运行前请先检查来源。
npx skills add https://github.com/ffsshhttiikk/opencode-agents-skills --skill etl-pipelines命令会保持在同一行。复制前请横向滚动并检查完整内容。
想先保存到本地?可下载 SkillsMP 当前能够提供的文件。
基于 SOC 职业分类
正在显示 SKILL.md
| name | etl-pipelines |
| description | ETL pipeline design and implementation |
| license | MIT |
| compatibility | opencode |
| metadata | {"audience":"data-engineers","category":"data-science"} |
Use me when:
# Apache Airflow DAG
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
"owner": "data_team",
"depends_on_past": False,
"retries": 3,
"retry_delay": timedelta(minutes=5)
}
with DAG(
"etl_pipeline",
start_date=datetime(2024, 1, 1),
schedule_interval="@daily",
default_args=default_args
) as dag:
extract = PythonOperator(
task_id="extract",
python_callable=extract_data
)
transform = PythonOperator(
task_id="transform",
python_callable=transform_data,
dependencies=[extract]
)
load = PythonOperator(
task_id="load",
python_callable=load_data,
dependencies=[transform]
)
from pyspark.sql import SparkSession
from pyspark.sql.functions import window, col
spark = SparkSession.builder \
.appName("StreamingETL") \
.config("spark.sql.streaming.checkpointLocation", "/checkpoints") \
.getOrCreate()
# Read streaming data
stream = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "broker:9092") \
.option("subscribe", "events") \
.load()
# Transform streaming
enriched = stream \
.select(
col("value").cast("string").alias("json")
) \
.withColumn("data", F.from_json("json", schema)) \
.select("data.*")
# Write streaming
query = enriched \
.writeStream \
.format("delta") \
.option("checkpointLocation", "/checkpoints") \
.trigger(processingTime="1 minute") \
.start("/data/lakehouse/tables/events")