Fabric Notebook读取Parquet写入Delta表时遇权限及实例错误
Monday API数据写入Delta表问题排查与解决
刚接触Python和公司的Fabric,从Monday API获取数据并写入Parquet文件流程正常,但写入Delta表时遇到两类错误:整合脚本报Java实例错误,拆分脚本报未授权错误,以下是问题分析及解决方法:
一、整合脚本:Java实例错误
代码
import subprocess import json import requests import pandas as pd import pyarrow as pa import pyarrow.parquet as pq from pyspark.sql import SparkSession access_token = "TOKEN" # Define your API endpoint and headers url = "https://api.monday.com/v2" headers = { "Authorization": access_token, "Content-Type": "application/json" } # Define your GraphQL query query = """ query { users { created_at email account { name id } } } """ # Send the request response = requests.post(url, json={'query': query}, headers=headers) # Check for errors if response.status_code == 200: data = response.json() users = data['data']['users'] # Convert the data to a DataFrame df = pd.json_normalize(users) parquet_table_name = "MondayUsers.parquet" parquet_file_path = "abfss://Files/Monday/" df.to_parquet(parquet_file_path + parquet_table_name, engine='pyarrow', index=False) print("Data written to Parquet file") # Initialize Spark session with Delta Lake configurations spark = SparkSession.builder \ .appName("FabricNotebook") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() parquet_file_path = parquet_file_path + parquet_table_name # Read the Parquet file into a Spark DataFrame try: from pyspark.sql import SparkSession spark = SparkSession.builder.appName("FabricNotebook").getOrCreate() df = spark.read.parquet(parquet_file_path) delta_table_name = "monday.Users" df.write.mode("overwrite").option("overwriteSchema", "true").format("delta").saveAsTable(delta_table_name) except Exception as e: print(f"Error reading Parquet file or writing to Delta table: {e}") else: print(f"Query failed with status code {response.status_code}")
错误信息(翻译后)
读取Parquet文件或写入Delta表时出错:调用o6842.toString时发生错误。跟踪信息: java.lang.IllegalArgumentException: 对象不是声明类的实例 at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.base/java.lang.Thread.run(Thread.java:829)
问题原因及解决
错误源于重复创建SparkSession导致上下文冲突:脚本先初始化了带Delta Lake配置的SparkSession,后续try块中又重新创建了未包含Delta配置的SparkSession,导致Java反射调用时出现实例不匹配。
修改方案:删除try块中重复创建SparkSession的代码,直接复用之前初始化好的Spark实例:
# 保留初始的SparkSession配置 spark = SparkSession.builder \ .appName("FabricNotebook") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() parquet_file_path = parquet_file_path + parquet_table_name try: # 移除这里重复创建SparkSession的代码 df = spark.read.parquet(parquet_file_path) delta_table_name = "monday.Users" df.write.mode("overwrite").option("overwriteSchema", "true").format("delta").saveAsTable(delta_table_name) except Exception as e: print(f"Error reading Parquet file or writing to Delta table: {e}")
二、拆分脚本:未授权错误
代码
import pandas as pd parquet_table_name = "MondayUsers.parquet" parquet_file_path = "abfss://Files/Monday/" parquet_file = parquet_file_path + parquet_table_name parquet_df = pd.read_parquet(parquet_file) spark_df = spark.createDataFrame(parquet_df) print(spark_df) delta_table_name = "monday.Users" spark_df.write.format("delta").mode("overwrite").saveAsTable(delta_table_name)
错误信息
未授权(unauthorized)
问题原因及解决
错误是由于Spark会话缺少访问Fabric资源的权限,或是通过pandas读取Parquet时未传递正确的身份凭证。
解决方法:
- 使用Fabric内置Spark上下文:Fabric Notebook中默认提供已授权的Spark会话,无需手动配置,直接使用即可
- 避免pandas中转:直接用Spark读取Parquet文件,减少权限传递问题
- 检查权限配置:确保当前用户拥有目标Lakehouse和
monday.Users表的读写权限
优化后的拆分脚本:
from pyspark.sql import SparkSession # 使用Fabric内置的Spark会话,自动继承权限 spark = SparkSession.builder.getOrCreate() parquet_table_name = "MondayUsers.parquet" parquet_file_path = "abfss://Files/Monday/" parquet_file = parquet_file_path + parquet_table_name # 直接用Spark读取Parquet spark_df = spark.read.parquet(parquet_file) delta_table_name = "monday.Users" spark_df.write.format("delta").mode("overwrite").saveAsTable(delta_table_name)
内容的提问来源于stack exchange,提问作者NutellaKing
相关产品推荐
相关产品推荐

