Databricks Runtime13.3中sc.parallelize兼容Unity Catalog问题及替代方案
Databricks Runtime版本升级后大体积JSON响应的序列化问题
原代码(Runtime 10中运行正常)
request_url = "https://.com/?fct=Get" # 修正原代码缺失的引号 response_task = requests.get(url=request_url, headers=headers) db1 = spark.sparkContext.parallelize([response_task.text]) df2 = spark.read.json(db1)
Runtime 13.3中运行报错
An error occurred while calling o398.createRDDFromTrustedPath. Trace: py4j.security.Py4JSecurityException: Method public org.apache.spark.api.java.JavaRDD org.apache.spark.sql.SparkSession.createRDDFromTrustedPath(java.lang.String,int) is not whitelisted on class class org.apache.spark.sql.SparkSession at py4j.security.WhitelistingPy4JSecurityManager.checkCall(WhitelistingPy4JSecurityManager.java:473) at py4j.Gateway.invoke(Gateway.java:305) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:195) at py4j.ClientServerConnection.run(ClientServerConnection.java:115) at java.lang.Thread.run(Thread.java:750)
尝试修改后的代码及报错
response_task = requests.get(url=request_url_gettasks, headers=headers) response_text = response_task.text db1 = spark.read.json([response_text])
报错信息:
com.databricks.rpc.UnexpectedHttpException: Got invalid response: 502
当前待大体积数据验证的可行代码
response_task = requests.get(url=request_url_gettasks, headers=headers) response_text = response_task.text task_data = json.loads(response_text)["results"] schema = StructType([ StructField("PROJECT_ID", IntegerType()), StructField("PROJECT_NAME", StringType()), StructField("PROJECT_NUMBER", StringType()), StructField("PROJECT_TYPE_NAME", StringType()), # 其他字段省略 ]) df2 = spark.createDataFrame(task_data, schema)
针对大体积response_task.text的其他序列化方式
写入DBFS临时文件后读取:适合超大规模数据,避免内存溢出
import tempfile # 将JSON文本写入DBFS临时文件 with tempfile.NamedTemporaryFile(mode='w', delete=False, dir='/dbfs/tmp/') as f: f.write(response_task.text) temp_path = f.name.replace('/dbfs', 'dbfs:/') # Spark读取临时文件中的JSON df = spark.read.json(temp_path) # 清理临时文件 dbutils.fs.rm(temp_path)分块解析JSON并分批处理:若响应为JSON数组,可拆分批次创建DataFrame
import json # 解析完整JSON数据 data = json.loads(response_task.text)["results"] batch_size = 10000 # 分批次处理数据 for i in range(0, len(data), batch_size): batch_data = data[i:i+batch_size] df_batch = spark.createDataFrame(batch_data, schema) # 按需写入表或进行后续处理 df_batch.write.mode("append").saveAsTable("target_table")Spark直接读取HTTP数据源:跳过本地请求,由集群分布式处理,适配大体积数据场景
df = spark.read.format("json").options(headers=str(headers)).load(request_url)注:需确保目标URL允许Spark集群访问,且请求头格式符合Spark要求。
内容的提问来源于stack exchange,提问作者Ahmed Bedhiaf
相关产品推荐
相关产品推荐

