如何用Java/Python将Delta Table/Synapse Table数据暴露为REST API并求示例
Python示例(FastAPI + PySpark)
用FastAPI快速搭建REST API,通过PySpark读取Delta Table或Synapse Table数据。
依赖安装
pip install fastapi uvicorn pyspark delta-spark
核心代码实现
from fastapi import FastAPI from pyspark.sql import SparkSession # 初始化支持Delta Lake的Spark Session spark = SparkSession.builder \ .appName("TableExposureAPI") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() app = FastAPI() # 暴露Delta Table数据接口 @app.get("/delta/{table_path}") def fetch_delta_data(table_path: str): try: # 读取Delta表,table_path可以是本地路径或云存储路径 df = spark.read.format("delta").load(table_path) # 返回前100条数据(避免大数据量返回) return {"status": "success", "data": df.limit(100).toJSON().collect()} except Exception as e: return {"status": "failed", "error_msg": str(e)} # 暴露Synapse Table数据接口 @app.get("/synapse/{table_name}") def fetch_synapse_data(table_name: str): try: # 通过JDBC连接Synapse,替换为你的实际连接参数 df = spark.read \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<你的Synapse工作区>.sql.azuresynapse.net:1433;database=<数据库名>") \ .option("dbtable", table_name) \ .option("user", "<用户名>") \ .option("password", "<密码>") \ .load() return {"status": "success", "data": df.limit(100).toJSON().collect()} except Exception as e: return {"status": "failed", "error_msg": str(e)} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
使用说明
- 替换Synapse连接参数中的占位符为你的实际信息
- 启动服务后,访问
http://localhost:8000/docs可查看自动生成的API文档
Java示例(Spring Boot + Spark Java)
基于Spring Boot构建REST API,通过Spark Java读取Delta或Synapse Table数据。
Maven依赖(pom.xml)
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>3.4.1</version> </dependency> <dependency> <groupId>io.delta</groupId> <artifactId>delta-core_2.12</artifactId> <version>2.4.0</version> </dependency> <dependency> <groupId>com.microsoft.sqlserver</groupId> <artifactId>mssql-jdbc</artifactId> <version>12.4.1.jre11</version> </dependency> </dependencies>
核心代码实现
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RestController; import java.util.List; import java.util.Map; import java.util.stream.Collectors; @SpringBootApplication @RestController public class TableApiService { // 初始化Spark Session private static final SparkSession spark = SparkSession.builder() .appName("TableApiService") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .master("local[*]") // 生产环境需替换为集群模式 .getOrCreate(); public static void main(String[] args) { SpringApplication.run(TableApiService.class, args); } // Delta Table数据接口 @GetMapping("/delta/{tablePath}") public Map<String, Object> getDeltaData(@PathVariable String tablePath) { try { Dataset<Row> df = spark.read().format("delta").load(tablePath); List<Map<String, Object>> dataList = df.limit(100).collectAsList().stream() .map(row -> row.getValuesMap(row.schema().fieldNames())) .collect(Collectors.toList()); return Map.of("status", "success", "data", dataList); } catch (Exception e) { return Map.of("status", "failed", "error_msg", e.getMessage()); } } // Synapse Table数据接口 @GetMapping("/synapse/{tableName}") public Map<String, Object> getSynapseData(@PathVariable String tableName) { try { Dataset<Row> df = spark.read() .format("jdbc") .option("url", "jdbc:sqlserver://<你的Synapse工作区>.sql.azuresynapse.net:1433;database=<数据库名>") .option("dbtable", tableName) .option("user", "<用户名>") .option("password", "<密码>") .load(); List<Map<String, Object>> dataList = df.limit(100).collectAsList().stream() .map(row -> row.getValuesMap(row.schema().fieldNames())) .collect(Collectors.toList()); return Map.of("status", "success", "data", dataList); } catch (Exception e) { return Map.of("status", "failed", "error_msg", e.getMessage()); } } }
使用说明
- 替换Synapse连接参数中的占位符为实际信息
- 生产环境需修改Spark的
master配置为对应集群模式 - 启动服务后,直接访问
http://localhost:8080/delta/<表路径>或/synapse/<表名>获取数据
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

