You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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时未传递正确的身份凭证。

解决方法:

  1. 使用Fabric内置Spark上下文:Fabric Notebook中默认提供已授权的Spark会话,无需手动配置,直接使用即可
  2. 避免pandas中转:直接用Spark读取Parquet文件,减少权限传递问题
  3. 检查权限配置:确保当前用户拥有目标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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 15:47:15