本篇文章Building Robust PySpark Code for Consistency, Optimization, and Flexibility requesting from API适合数据工程师学习如何构建稳定的PySpark管道。文章的亮点在于采用了动态配置连接、批处理API请求和分布式处理等方法,确保数据一致性和优化性能。



PySpark Pipeline

在快节奏的数据工程世界中,维护稳定、高效且适应性强的代码对于无缝处理大型数据集至关重要。本文将详细介绍我如何开发一个健壮的 PySpark 数据管道,以确保数据一致性、优化和灵活性。

1 挑战

在 Databricks 和 Snowflake 等分布式数据环境中工作时,我需要一个可扩展且具有错误恢复能力的管道,以高效处理大规模业务数据。该管道必须:

  • 确保数据一致性:在数据获取和处理过程中保持一致性。
  • 优化性能:优化 API 调用和数据库查询以提高性能。
  • 保持灵活性:适应新的数据结构和业务需求。
  • 优雅地处理错误:最大程度地减少停机时间,确保数据流畅。
  • 提高可维护性:通过模块化和可重用代码来提升可维护性。

2 设计可扩展解决方案

为了实现这些目标,我使用 PySpark、REST API 调用和 Snowflake 集成相结合的方式构建了数据管道。以下是我的方法中的关键组件:

2.1 可配置的 Snowflake 连接

我没有硬编码凭据,而是利用安全的凭据存储来动态管理身份验证。这确保了安全性,同时使得修改连接设置而无需更改代码库变得容易。

options = {  
    "sfUrl": "<snowflake_url>",  
    "sfUser": dbutils.secrets.get("scope", "dbuser"),  
    "sfPassword": dbutils.secrets.get("scope", "dbpasswd"),  
    "sfDatabase": "DATABASE_NAME",  
    "sfSchema": "SCHEMA_NAME",  
    "sfWarehouse": "WAREHOUSE_NAME"  
}

2.2 高效的 API 请求处理

由于数据管道与外部 API 交互,我实现了带有重试逻辑的批量处理,以优雅地处理请求失败。

def process_batch(ids, selection, retries=5, wait_time=3):  
    query = json.dumps({"WHERE": [{"ID": ids}], "SELECT": selection})  
    for attempt in range(retries):  
        try:  
            response = requests.post("<api_url>", headers=headers, data=query, timeout=30)  
            if response.status_code == 200:  
                return response.json().get('Data', [])  
            elif response.status_code >= 500:  
                time.sleep(wait_time)  
        except (ReadTimeout, ConnectionError):  
            time.sleep(wait_time)  
    return []

2.3 使用 PySpark 进行分布式处理

我没有迭代单个记录,而是使用 PySpark 的 RDD 转换将处理分布到多个节点。

def process_partition(partition):  
    results = []  
    for batch in partition:  
        results.extend(process_batch(batch, selection))  
    return iter(results)

id_df = spark.createDataFrame([(id,) for id in id_list], ["id"])  
result_rdd = id_df.rdd.mapPartitions(process_partition)  
data = result_rdd.collect()

2.4 Snowflake 中的查询优化

为了最大限度地减少对 Snowflake 性能的影响,我设计了优化的查询,使用了过滤、索引和 DISTINCT 选择。

query = """  
    (SELECT DISTINCT ID FROM SCHEMA.TABLE  
     WHERE CONDITION_1 AND CONDITION_2)  
"""  
sparkDF = spark.read.format("snowflake").options(**options).option("query", query).load()

2.5 动态数据选择和适应性

为了确保灵活性,我根据业务需求动态构建了 API 选择标准。

def generate_selection(index_list):  
    return [  
        {"ID": {"AS": "ID"}},  
        {"NAME": {"AS": "NAME"}},  
        {"CURRENCY": {"INDEX": index_list, "AS": "CURRENCY"}},  
    ]

2.6 错误处理和日志记录

为了进一步增强稳定性,我引入了结构化日志记录和异常处理,以主动监控问题。

import logging

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')  
logger = logging.getLogger(__name__)

def safe_query_execution(query):  
    try:  
        return spark.read.format("snowflake").options(**options).option("query", query).load()  
    except Exception as e:  
        logger.error(f"Query execution failed: {str(e)}")  
        return None

3 成果

通过遵循这些原则,我构建了一个 PySpark 数据管道,它:

  • ✅ 通过结构化 API 查询和 Snowflake 约束保证了数据一致性
  • ✅ 通过批量请求和利用分布式处理优化了性能。
  • ✅ 动态适应不断变化的业务需求,无需进行重大代码修改。
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐