如何在Azure中以Parquet/Delta Parquet格式存储百万级数据及适配CRUD API?
将百万级CRUD数据存储为Parquet/Delta Parquet的方案及Azure相关服务
一、Parquet/Delta Parquet存储实现方法
1. 批量数据写入
针对现有批量数据的格式转换,可通过编程语言库或分布式计算工具实现:
- Python + Pandas/PyArrow:适合中小批量数据,操作轻量化
import pandas as pd # 示例:从API返回的JSON数据构建DataFrame api_data_df = pd.read_json("api_response.json") # 写入压缩Parquet文件 api_data_df.to_parquet("data/user_data.parquet", engine="pyarrow", compression="snappy")
- Spark + Delta Lake:适配百万级以上大规模数据,支持分布式处理
from pyspark.sql import SparkSession # 初始化带Delta Lake支持的Spark会话 spark = SparkSession.builder \ .appName("BatchDeltaWrite") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() # 读取API数据转为Spark DataFrame api_data_df = spark.read.json("api_response.json") # 写入Delta表(支持后续增删改操作) api_data_df.write.format("delta").mode("overwrite").save("abfss://container@storageaccount.dfs.core.windows.net/delta_tables/user_data")
2. 增量CRUD操作
Parquet本身是静态列存格式,无法高效支持单条数据的增删改,Delta Parquet(Delta Lake)是CRUD场景的最优选择,它提供ACID事务、版本控制、增量更新能力:
- 单条/批量增删改:通过Delta Lake SDK实现
from delta.tables import DeltaTable from pyspark.sql.functions import current_timestamp delta_table = DeltaTable.forPath(spark, "abfss://container@storageaccount.dfs.core.windows.net/delta_tables/user_data") # 更新数据 delta_table.update( condition="user_id = 'U1001'", set={"user_name": "UpdatedName", "last_update": current_timestamp()} ) # 插入新数据(避免重复) new_data_df = spark.createDataFrame([("U1002", "NewUser", current_timestamp())], ["user_id", "user_name", "create_time"]) delta_table.alias("t").merge( new_data_df.alias("s"), "t.user_id = s.user_id" ).whenNotMatchedInsertAll().execute() # 删除数据 delta_table.delete(condition="user_id = 'U1000'")
- 实时CRUD处理:如果API是高并发实时请求,可通过Flink结合Delta Lake做流式写入,将API请求转为数据流后写入Delta表,保障数据一致性。
二、Azure中支持Parquet/Delta Parquet的存储及API服务
Azure提供多类服务覆盖Parquet/Delta格式的存储、计算与查询需求:
- Azure Data Lake Storage Gen2(ADLS Gen2):作为底层存储介质,专门用于存储Parquet/Delta文件,支持通过Azure Storage SDK、REST API直接操作文件,也可配合其他计算服务实现结构化查询。
- Azure Databricks:原生支持Delta Lake,提供类SQL的查询API以及Python/Scala SDK,完美适配CRUD场景。可通过SQL直接操作Delta表:
-- 创建Delta表 CREATE TABLE user_data ( user_id STRING, user_name STRING, create_time TIMESTAMP ) USING DELTA LOCATION 'abfss://container@storageaccount.dfs.core.windows.net/delta_tables/user_data'; -- 执行更新操作 UPDATE user_data SET user_name = 'NewName' WHERE user_id = 'U1001';
- Azure Synapse Analytics:
- Serverless SQL池:可直接查询ADLS Gen2中的Parquet文件,支持T-SQL语法;
- Spark池:完整支持Delta Lake的CRUD操作,可与Synapse的ETL、分析组件联动;
- 支持创建外部表映射Parquet文件,实现类数据库的查询体验。
- Azure Cosmos DB(Synapse Link):若已有Cosmos DB数据,可通过Synapse Link将数据以Parquet格式导出到ADLS Gen2,或直接查询Parquet格式的分析存储,但原生CRUD仍以Cosmos DB的NoSQL API为主,并非专门的Parquet存储引擎。
注意:如果你的CRUD场景以高频单条操作为主,建议先将数据写入行存数据库(如Azure SQL Database),再通过CDC(变更数据捕获)或定时任务同步到Delta/Parquet存储用于分析;如果是批量CRUD或分析为核心场景,直接用Delta Lake + Azure Databricks/Synapse是更高效的方案。
内容的提问来源于stack exchange,提问作者Jerald Baker
相关产品推荐
相关产品推荐

