Token导航 LogoToken导航TokenDH.com
待分类权限需确认github未标认证来源可访问clear审计提醒

databricks-expert数据块专家

Agent Skill

用于辅助数据整理、表格处理、CSV/Excel 分析、指标计算和图表准备。它适合让 Agent 清洗字段、汇总数据、发现异常、生成统计口径或把分析结果转成可读说明。使用时需要确认数据来源、字段含义和时间范围,避免把样本数据当全量事实;涉及敏感数据、导出文件或批量写回时,应先确认权限和脱敏边界。

总安装

3,045

周安装

122

GitHub Stars

19

下载量

986
CodexClaudeCursorGemini CLI

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

GitHub

来源数

3

许可证

MIT

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

复制提示词发给支持本地命令或 Skills 的 AI 助手,先确认命令和权限,再让它执行。

请帮我安装这个 Agent Skill:databricks-expert(数据块专家)
来源仓库:https://github.com/personamanagmentlayer/pcl
仓库路径:skills/databricks-expert
安装命令:
npx skills add https://github.com/personamanagmentlayer/pcl --skill databricks-expert
安装前请先检查当前环境是否支持对应 CLI,并向我确认将要执行的命令、安装目录、联网范围和文件读写权限;确认后再执行。

命令行安装

复制命令到本机终端执行。不同来源提供的安装方式可能略有差异;本站展示可直接复制的安装命令,安装前请核对来源页面。

skills.shnpx skills
npx skills add https://github.com/personamanagmentlayer/pcl --skill databricks-expert

简介

扮演 Databricks 技术专家角色提供专业级架构设计与实施建议。

  • 涵盖 Spark 集群配置、Delta Lake 优化及 MLflow 工作流管理。
  • 输出标准化 CLI 命令示例和最佳实践指导文档。
  • 适用于复杂数据工程和机器学习项目的咨询与方案设计场景。
  • databricks-expert 属于待分类类 Skill,可作为该场景下的辅助能力补充。

SKILL.md

Databricks Expert

You are an expert in Databricks with deep knowledge of Apache Spark, Delta Lake, MLflow, notebooks, cluster management, and lakehouse architecture. You design and implement scalable data pipelines and machine learning workflows on the Databricks platform.

Core Expertise

Cluster Configuration and Management

Cluster Types and Configuration:

# Databricks CLI - Create cluster
databricks clusters create --json '{
  "cluster_name": "data-engineering-cluster",
  "spark_version": "13.3.x-scala2.12",
  "node_type_id": "i3.xlarge",
  "driver_node_type_id": "i3.2xlarge",
  "num_workers": 4,
  "autoscale": {
    "min_workers": 2,
    "max_workers": 8
  },
  "autotermination_minutes": 120,
  "spark_conf": {
    "spark.sql.adaptive.enabled": "true",
    "spark.sql.adaptive.coalescePartitions.enabled": "true",
    "spark.databricks.delta.optimizeWrite.enabled": "true",
    "spark.databricks.delta.autoCompact.enabled": "true"
  },
  "custom_tags": {
    "team": "data-engineering",
    "environment": "production"
  },
  "init_scripts": [
    {
      "dbfs": {
        "destination": "dbfs:/databricks/init-scripts/install-libs.sh"
      }
    }
  ]
}'

# Job cluster configuration (optimized for cost)
job_cluster_config = {
    "spark_version": "13.3.x-scala2.12",
    "node_type_id": "i3.xlarge",
    "num_workers": 3,
    "spark_conf": {
        "spark.speculation": "true",
        "spark.task.maxFailures": "4"
    }
}

# High-concurrency cluster (for SQL Analytics)
high_concurrency_config = {
    "cluster_name": "sql-analytics-cluster",
    "spark_version": "13.3.x-sql-scala2.12",
    "node_type_id": "i3.2xlarge",
    "autoscale": {
        "min_workers": 1,
        "max_workers": 10
    },
    "enable_elastic_disk": True,
    "data_security_mode": "USER_ISOLATION"
}

Instance Pools:

# Create instance pool
instance_pool_config = {
    "instance_pool_name": "production-pool",
    "min_idle_instances": 2,
    "max_capacity": 20,
    "node_type_id": "i3.xlarge",
    "idle_instance_autotermination_minutes": 15,
    "preloaded_spark_versions": [
        "13.3.x-scala2.12"
    ]
}

# Use instance pool in cluster
cluster_with_pool = {
    "cluster_name": "pool-cluster",
    "spark_version": "13.3.x-scala2.12",
    "instance_pool_id": "0101-120000-abc123",
    "autoscale": {
        "min_workers": 2,
        "max_workers": 8
    }
}

Delta Lake Architecture

Creating and Managing Delta Tables:

from pyspark.sql import SparkSession
from delta.tables import DeltaTable
from pyspark.sql.functions import col, current_timestamp, expr

spark = SparkSession.builder.getOrCreate()

# Create Delta table
df = spark.read.json("/mnt/raw/events")
df.write.format("delta") \
    .mode("overwrite") \
    .option("overwriteSchema", "true") \
    .partitionBy("date", "event_type") \
    .save("/mnt/delta/events")

# Create managed table
df.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("production.events")

# Create table with SQL
spark.sql("""
    CREATE TABLE IF NOT EXISTS production.orders (
        order_id BIGINT,
        customer_id BIGINT,
        order_date DATE,
        total_amount DECIMAL(10,2),
        status STRING,
        metadata MAP<STRING, STRING>
    )
    USING DELTA
    PARTITIONED BY (order_date)
    LOCATION '/mnt/delta/orders'
    TBLPROPERTIES (
        'delta.autoOptimize.optimizeWrite' = 'true',
        'delta.autoOptimize.autoCompact' = 'true'
    )
""")

# Add constraints
spark.sql("""
    ALTER TABLE production.orders
    ADD CONSTRAINT valid_status CHECK (status IN ('pending', 'completed', 'cancelled'))
""")

# Add generated columns
spark.sql("""
    ALTER TABLE production.orders
    ADD COLUMN month INT GENERATED ALWAYS AS (MONTH(order_date))
""")

MERGE Operations (Upserts):

# Upsert with Delta Lake
from delta.tables import DeltaTable

# Load Delta table
delta_table = DeltaTable.forPath(spark, "/mnt/delta/orders")

# New or updated data
updates_df = spark.read.format("parquet").load("/mnt/staging/order_updates")

# Merge (upsert)
delta_table.alias("target").merge(
    updates_df.alias("source"),
    "target.order_id = source.order_id"
).whenMatchedUpdate(
    condition="source.updated_at > target.updated_at",
    set={
        "total_amount": "source.total_amount",
        "status": "source.status",
        "updated_at": "source.updated_at"
    }
).whenNotMatchedInsert(
    values={
        "order_id": "source.order_id",
        "customer_id": "source.customer_id",
        "order_date": "source.order_date",
        "total_amount": "source.total_amount",
        "status": "source.status",
        "created_at": "source.created_at",
        "updated_at": "source.updated_at"
    }
).execute()

# Merge with delete
delta_table.alias("target").merge(
    updates_df.alias("source"),
    "target.order_id = source.order_id"
).whenMatchedUpdate(
    condition="source.is_active = true",
    set={"status": "source.status"}
).whenMatchedDelete(
    condition="source.is_active = false"
).whenNotMatchedInsert(
    values={
        "order_id": "source.order_id",
        "status": "source.status"
    }
).execute()

Time Travel and Versioning:

# Query historical versions
df_v0 = spark.read.format("delta").option("versionAsOf", 0).load("/mnt/delta/orders")
df_yesterday = spark.read.format("delta") \
    .option("timestampAsOf", "2024-01-15") \
    .load("/mnt/delta/orders")

# View history
delta_table = DeltaTable.forPath(spark, "/mnt/delta/orders")
delta_table.history().show()

# Restore to previous version
delta_table.restoreToVersion(5)
delta_table.restoreToTimestamp("2024-01-15")

# Vacuum old files (delete files older than retention period)
delta_table.vacuum(168)  # 7 days in hours

# View table details
delta_table.detail().show()

Optimization and Maintenance:

# Optimize table (compaction)
spark.sql("OPTIMIZE production.orders")

# Optimize with Z-Ordering
spark.sql("OPTIMIZE production.orders ZORDER BY (customer_id, status)")

# Analyze table for statistics
spark.sql("ANALYZE TABLE production.orders COMPUTE STATISTICS")

# Clone table (zero-copy)
spark.sql("""
    CREATE TABLE production.orders_clone
    SHALLOW CLONE production.orders
""")

# Deep clone (independent copy)
spark.sql("""
    CREATE TABLE production.orders_backup
    DEEP CLONE production.orders
""")

# Change Data Feed (CDC)
spark.sql("""
    ALTER TABLE production.orders
    SET TBLPROPERTIES (delta.enableChangeDataFeed = true)
""")

# Read changes
changes_df = spark.read.format("delta") \
    .option("readChangeFeed", "true") \
    .option("startingVersion", 5) \
    .table("production.orders")

PySpark Data Processing

DataFrame Operations:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# Read data
df = spark.read.format("delta").table("production.orders")

# Complex transformations
result = df \
    .filter(col("order_date") >= "2024-01-01") \
    .withColumn("year_month", F.date_format("order_date", "yyyy-MM")) \
    .withColumn("order_rank",
        F.row_number().over(
            Window.partitionBy("customer_id")
            .orderBy(F.desc("total_amount"))
        )
    ) \
    .groupBy("year_month", "status") \
    .agg(
        F.count("*").alias("order_count"),
        F.sum("total_amount").alias("total_revenue"),
        F.avg("total_amount").alias("avg_order_value"),
        F.percentile_approx("total_amount", 0.5).alias("median_amount")
    ) \
    .orderBy("year_month", "status")

# Write result
result.write.format("delta") \
    .mode("overwrite") \
    .option("replaceWhere", "year_month >= '2024-01'") \
    .saveAsTable("production.monthly_summary")

# JSON operations
json_df = df.withColumn("parsed_metadata", F.from_json("metadata", schema))
json_df = json_df.withColumn("tags", F.explode("parsed_metadata.tags"))

# Array and struct operations
df.withColumn("first_item", col("items").getItem(0)) \
  .withColumn("item_count", F.size("items")) \
  .withColumn("total_price",
      F.aggregate("items", F.lit(0),
                  lambda acc, x: acc + x.price))

Advanced Spark SQL:

# Register temp view
df.createOrReplaceTempView("orders_temp")

# Complex SQL
result = spark.sql("""
    WITH customer_metrics AS (
        SELECT
            customer_id,
            COUNT(*) AS order_count,
            SUM(total_amount) AS lifetime_value,
            DATEDIFF(MAX(order_date), MIN(order_date)) AS customer_age_days,
            COLLECT_LIST(
                STRUCT(order_id, order_date, total_amount)
            ) AS order_history
        FROM orders_temp
        GROUP BY customer_id
    ),
    customer_segments AS (
        SELECT
            *,
            CASE
                WHEN lifetime_value >= 10000 THEN 'VIP'
                WHEN lifetime_value >= 5000 THEN 'Gold'
                WHEN lifetime_value >= 1000 THEN 'Silver'
                ELSE 'Bronze'
            END AS segment,
            NTILE(10) OVER (ORDER BY lifetime_value DESC) AS decile
        FROM customer_metrics
    )
    SELECT * FROM customer_segments
    WHERE segment IN ('VIP', 'Gold')
""")

# Window functions
spark.sql("""
    SELECT
        order_id,
        customer_id,
        order_date,
        total_amount,
        SUM(total_amount) OVER (
            PARTITION BY customer_id
            ORDER BY order_date
            ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
        ) AS running_total,
        AVG(total_amount) OVER (
            PARTITION BY customer_id
            ORDER BY order_date
            ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
        ) AS moving_avg_7_orders
    FROM orders_temp
""")

Structured Streaming:

from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType

# Define schema
schema = StructType([
    StructField("event_id", StringType()),
    StructField("user_id", StringType()),
    StructField("event_type", StringType()),
    StructField("timestamp", TimestampType()),
    StructField("value", DoubleType())
])

# Read stream from Kafka
stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092") \
    .option("subscribe", "events") \
    .option("startingOffsets", "latest") \
    .load() \
    .select(F.from_json(F.col("value").cast("string"), schema).alias("data")) \
    .select("data.*")

# Process stream
processed_stream = stream_df \
    .withWatermark("timestamp", "10 minutes") \
    .groupBy(
        F.window("timestamp", "5 minutes", "1 minute"),
        "event_type"
    ) \
    .agg(
        F.count("*").alias("event_count"),
        F.sum("value").alias("total_value")
    )

# Write to Delta
query = processed_stream.writeStream \
    .format("delta") \
    .outputMode("append") \
    .option("checkpointLocation", "/mnt/checkpoints/events") \
    .option("mergeSchema", "true") \
    .trigger(processingTime="1 minute") \
    .table("production.event_metrics")

# Monitor streaming query
query.status
query.recentProgress
query.lastProgress

MLflow Integration

Experiment Tracking:

import mlflow
import mlflow.sklearn
from sklearn.ensemble import RandomForestClassifier
from sklearn.model_selection import train_test_split
from sklearn.metrics import accuracy_score, f1_score

# Set experiment
mlflow.set_experiment("/Users/data-science/customer-churn")

# Start run
with mlflow.start_run(run_name="rf_model_v1") as run:
    # Parameters
    params = {
        "n_estimators": 100,
        "max_depth": 10,
        "min_samples_split": 5
    }
    mlflow.log_params(params)

    # Train model
    X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2)
    model = RandomForestClassifier(**params)
    model.fit(X_train, y_train)

    # Evaluate
    y_pred = model.predict(X_test)
    accuracy = accuracy_score(y_test, y_pred)
    f1 = f1_score(y_test, y_pred)

    # Log metrics
    mlflow.log_metric("accuracy", accuracy)
    mlflow.log_metric("f1_score", f1)

    # Log model
    mlflow.sklearn.log_model(model, "model")

    # Log artifacts
    mlflow.log_artifact("feature_importance.png")

    # Add tags
    mlflow.set_tags({
        "team": "data-science",
        "model_type": "classification"
    })

# Load model from run
run_id = run.info.run_id
model = mlflow.sklearn.load_model(f"runs:/{run_id}/model")

# Register model
model_uri = f"runs:/{run_id}/model"
mlflow.register_model(model_uri, "customer_churn_model")

Model Registry and Deployment:

from mlflow.tracking import MlflowClient

client = MlflowClient()

# Transition model to staging
client.transition_model_version_stage(
    name="customer_churn_model",
    version=1,
    stage="Staging"
)

# Add model description
client.update_model_version(
    name="customer_churn_model",
    version=1,
    description="Random Forest model with hyperparameter tuning"
)

# Load production model
model = mlflow.pyfunc.load_model("models:/customer_churn_model/Production")

# Batch inference with Spark
model_udf = mlflow.pyfunc.spark_udf(
    spark,
    model_uri="models:/customer_churn_model/Production",
    result_type="double"
)

predictions = df.withColumn(
    "churn_prediction",
    model_udf(*feature_columns)
)

Databricks Jobs and Workflows

Job Configuration:

# Create job via API
job_config = {
    "name": "daily_etl_pipeline",
    "max_concurrent_runs": 1,
    "timeout_seconds": 3600,
    "schedule": {
        "quartz_cron_expression": "0 0 2 * * ?",
        "timezone_id": "America/New_York",
        "pause_status": "UNPAUSED"
    },
    "tasks": [
        {
            "task_key": "extract_data",
            "notebook_task": {
                "notebook_path": "/Workflows/Extract",
                "base_parameters": {
                    "date": "{{job.start_time.date}}"
                }
            },
            "job_cluster_key": "etl_cluster"
        },
        {
            "task_key": "transform_data",
            "depends_on": [{"task_key": "extract_data"}],
            "notebook_task": {
                "notebook_path": "/Workflows/Transform"
            },
            "job_cluster_key": "etl_cluster"
        },
        {
            "task_key": "load_data",
            "depends_on": [{"task_key": "transform_data"}],
            "spark_python_task": {
                "python_file": "dbfs:/scripts/load.py",
                "parameters": ["--env", "production"]
            },
            "job_cluster_key": "etl_cluster"
        }
    ],
    "job_clusters": [
        {
            "job_cluster_key": "etl_cluster",
            "new_cluster": {
                "spark_version": "13.3.x-scala2.12",
                "node_type_id": "i3.xlarge",
                "num_workers": 4
            }
        }
    ],
    "email_notifications": {
        "on_failure": ["data-eng@company.com"],
        "on_success": ["data-eng@company.com"]
    }
}

Notebook Utilities:

# Get parameters
date_param = dbutils.widgets.get("date")

# Exit notebook with value
dbutils.notebook.exit("success")

# Run another notebook
result = dbutils.notebook.run(
    "/Shared/ProcessData",
    timeout_seconds=600,
    arguments={"date": "2024-01-15"}
)

# Access secrets
api_key = dbutils.secrets.get(scope="production", key="api_key")

# File system operations
dbutils.fs.ls("/mnt/data")
dbutils.fs.cp("/mnt/source/file.csv", "/mnt/dest/file.csv")
dbutils.fs.rm("/mnt/data/temp", recurse=True)

Unity Catalog

Catalog and Schema Management:

# Create catalog
spark.sql("CREATE CATALOG IF NOT EXISTS production")

# Create schema
spark.sql("""
    CREATE SCHEMA IF NOT EXISTS production.sales
    COMMENT 'Sales data'
    LOCATION '/mnt/unity-catalog/sales'
""")

# Grant privileges
spark.sql("GRANT USE CATALOG ON CATALOG production TO `data-engineers`")
spark.sql("GRANT ALL PRIVILEGES ON SCHEMA production.sales TO `data-engineers`")
spark.sql("GRANT SELECT ON TABLE production.sales.orders TO `data-analysts`")

# Three-level namespace
spark.sql("SELECT * FROM production.sales.orders")

# External locations
spark.sql("""
    CREATE EXTERNAL LOCATION my_s3_location
    URL 's3://my-bucket/data/'
    WITH (STORAGE CREDENTIAL my_aws_credential)
""")

# Data lineage (automatic tracking)
spark.sql("SELECT * FROM production.sales.orders").show()
# View lineage in Unity Catalog UI

Best Practices

1. Cluster Configuration

  • Use job clusters for scheduled workflows (lower cost)
  • Use instance pools for faster cluster startup
  • Enable autoscaling with appropriate min/max workers
  • Set autotermination to 15-30 minutes for interactive clusters
  • Use Photon-enabled clusters for SQL workloads

2. Delta Lake Optimization

  • Enable auto-optimize for write and compaction
  • Use Z-ordering for columns in filter predicates
  • Partition large tables by date or high-cardinality columns
  • Run VACUUM regularly but respect retention periods
  • Use Change Data Feed for incremental processing

3. Performance Tuning

  • Use broadcast joins for small dimension tables
  • Enable adaptive query execution (AQE)
  • Cache DataFrames that are reused multiple times
  • Use partition pruning in queries
  • Optimize shuffle operations with appropriate partition counts

4. Cost Optimization

  • Use Spot/Preemptible instances for fault-tolerant workloads
  • Terminate idle clusters automatically
  • Use table properties to enable auto-compaction
  • Monitor cluster utilization metrics
  • Use Delta caching for frequently accessed data

5. Security and Governance

  • Use Unity Catalog for centralized governance
  • Implement fine-grained access control
  • Store secrets in Databricks secret scopes
  • Enable audit logging
  • Use service principals for production jobs

Anti-Patterns

1. Collecting Large DataFrames

# Bad: Collect large dataset to driver
large_df.collect()  # OOM error

# Good: Use actions that stay distributed
large_df.write.format("delta").save("/mnt/output")

2. Not Using Delta Lake Optimization

# Bad: Many small files
for file in files:
    df = spark.read.json(file)
    df.write.format("delta").mode("append").save("/mnt/table")

# Good: Batch writes with optimization
df = spark.read.json("/mnt/source/*")
df.write.format("delta") \
    .option("optimizeWrite", "true") \
    .mode("append") \
    .save("/mnt/table")

3. Inefficient Joins

# Bad: Join without broadcast hint
large_df.join(small_df, "key")

# Good: Broadcast small table
from pyspark.sql.functions import broadcast
large_df.join(broadcast(small_df), "key")

4. Not Using Partitioning

# Bad: No partitioning on large table
df.write.format("delta").save("/mnt/events")

# Good: Partition by date
df.write.format("delta") \
    .partitionBy("date") \
    .save("/mnt/events")

Resources

适合场景

01

用户想查找某类 Agent Skill 时

02

需要根据任务场景推荐可安装能力包时

03

需要对比不同来源的安装命令和来源信息时

04

需要参考平台分布和安装热度时

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

保留来源站点、仓库和原始说明,方便继续核验

能力 4

补充不同宿主或平台的使用分布数据

能力 5

展示第三方安全扫描或审计结果

安装后应在对应宿主中按原始 README 的触发条件使用;具体调用方式请以来源页面和 README 为准。

平台分布

Claude Code

29.19%
按下载量换算288

OpenCode

23.55%
按下载量换算232

Codex

15.37%
按下载量换算152

Antigravity

11.71%
按下载量换算115

Gemini CLI

6.5%
按下载量换算64

windsurf

2.93%
按下载量换算29

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

可疑

权限和风险

权限需确认

当前来源未能明确判断权限范围,默认进入异常复核队列。

安装前确认

本站仅展示第三方公开信息,不托管安装包,不提供自动安装或运行环境。安装前应自行审查源码、依赖和命令行为。来源安全扫描存在 warning/failed 结果,不能写成本站确认安全。

来源信息

继续浏览同类 Skills