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

Azure Synapse Pipeline执行报错errorCode:6002(K12参数失败)

K12参数下Azure Synapse PySpark管道Schema不匹配错误(errorCode:6002)

问题场景

  • 基于Microsoft Education Insights管道实现,在Azure Synapse Analytics工作区运行PySpark管道
  • 使用HED参数时执行正常,切换为K12参数后触发errorCode:6002错误
  • 错误核心:Delta目标表仅包含_c0至_c7共8个字段,但K12输入数据新增了_c8字段,导致写入Delta表时Schema不匹配

完整错误栈

{"errorCode": "6002","message": "---------------------------------------------------------------------------\nStreamingQueryException                   Traceback (most recent call last)\n/tmp/ipykernel_9514/162162073.py in <module>\n      1 metadata = oea.get_metadata_from_url('https://raw.githubusercontent.com/microsoft/OpenEduAnalytics/main/modules/module_catalog/Microsoft_Education_Insights/test_data/metadata.csv')\n----> 2 ingest_insights_dataset('stage1/Transactional/M365/v1.14')\n\n/tmp/ipykernel_9514/2788687738.py in ingest_insights_dataset(tables_source)\n     17                 logger.info('No test data')\n     18             else:\n---> 19                 oea.ingest('M365/v1.14/' + item, '_c0', options)\n     20         except AnalysisException as e:\n     21             # This means the table may have not been properly refined due to errors with the primary key not aligning with columns expected in the lookup table.\n\n/tmp/ipykernel_9514/1615070926.py in ingest(self, entity_path, primary_key, options)\n    474         if source_data_format == 'csv' and (not 'header' in options or options['header'] == None): options['header'] = True  # default to expecting a header in csv files\n    475 \n---> 476         number_of_new_inbound_rows = self.process(source_url, batch_func, options)\n    477         if number_of_new_inbound_rows > 0:\n    478             self.add_to_lake_db(ingested_path)\n\n/tmp/ipykernel_9514/1615070926.py in process(self, source_path, foreach_batch_function, options)\n    425         # for more info on append vs complete vs update modes for structured streaming: https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#basic-concepts\n    426         query = streaming_df.writeStream.format('delta').outputMode('append').trigger(once=True).option('checkpointLocation', self.to_url(source_path) + '/_checkpoints').foreachBatch(wrapped_function).start()\n---> 427         query.awaitTermination()   # block until query is terminated, with stop() or with error; A StreamingQueryException will be thrown if an exception occurs.\n    428         number_of_new_inbound_rows = query.lastProgress[\"numInputRows\"]\n    429         logger.info(f'Number of new inbound rows processed: {number_of_new_inbound_rows}')\n\n/opt/spark/python/lib/pyspark.zip/pyspark/sql/streaming.py in awaitTermination(self, timeout)\n     99             return self._jsq.awaitTermination(int(timeout * 1000))\n    100         else:\n---> 101             return self._jsq.awaitTermination()\n    102 \n    103     @property\n\n~/cluster-env/env/lib/python3.8/site-packages/py4j/java_gateway.py in __call__(self, *args)\n   1319 \n   1320         answer = self.gateway_client.send_command(command)\n-> 1321         return_value = get_return_value(\n   1322             answer, self.gateway_client, self.target_id, self.name)\n   1323 \n\n/opt/spark/python/lib/pyspark.zip/pyspark/sql/utils.py in deco(*a, **kw)\n    115                 # Hide where the exception came from that shows a non-Pythonic\n    116                 # JVM exception message.\n---> 117                 raise converted from None\n    118             else:\n    119                 raise\n\nStreamingQueryException: An exception was raised by the Python Proxy. Return Message: Traceback (most recent call last):\n  File \"/home/trusted-service-user/cluster-env/env/lib/python3.8/site-packages/py4j/java_gateway.py\", line 2463, in _call_proxy\n    return_value = getattr(self.pool[obj_id], method)(*params)\n  File \"/opt/spark/python/lib/pyspark.zip/pyspark/sql/utils.py\", line 196, in call\n    raise e\n  File \"/opt/spark/python/lib/pyspark.zip/pyspark/sql/utils.py\", line 193, in call\n    self.func(DataFrame(jdf, self.sql_ctx), batch_id)\n  File \"/tmp/ipykernel_9514/1615070926.py\", line 420, in wrapped_function\n    foreach_batch_function(df)\n  File \"/tmp/ipykernel_9514/1615070926.py\", line 464, in batch_func\n    def batch_func(df): self.overwrite(df, ingested_path, primary_key)\n  File \"/tmp/ipykernel_9514/1615070926.py\", line 389, in overwrite\n    df.write.format('delta').mode('overwrite').save(destination_url)  # https://docs.delta.io/latest/delta-batch.html#overwrite\n  File \"/opt/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py\", line 740, in save\n    self._jwrite.save(path)\n  File \"/home/trusted-service-user/cluster-env/env/lib/python3.8/site-packages/py4j/java_gateway.py\", line 1321, in __call__\n    return_value = get_return_value(\n  File \"/opt/spark/python/lib/pyspark.zip/pyspark/sql/utils.py\", line 117, in deco\n    raise converted from None\npyspark.sql.utils.AnalysisException: A schema mismatch detected when writing to the Delta table (Table ID: 5a0f37b0-5bf4-4906-9770-dd792df98a24).\nTo enable schema migration using DataFrameWriter or DataStreamWriter, please set:\n'.option(\"mergeSchema\", \"true\")'.\nFor other operations, set the session configuration\nspark.databricks.delta.schema.autoMerge.enabled to \"true\". See the documentation\nspecific to the operation for details.\n\nTable schema:\nroot\n-- _c0: string (nullable = true)\n-- _c1: string (nullable = true)\n-- _c2: string (nullable = true)\n-- _c3: string (nullable = true)\n-- _c4: string (nullable = true)\n-- _c5: string (nullable = true)\n-- _c6: string (nullable = true)\n-- _c7: string (nullable = true)\n\n\nData schema:\nroot\n-- _c0: string (nullable = true)\n-- _c1: string (nullable = true)\n-- _c2: string (nullable = true)\n-- _c3: string (nullable = true)\n-- _c4: string (nullable = true)\n-- _c5: string (nullable = true)\n-- _c6: string (nullable = true)\n-- _c7: string (nullable = true)\n-- _c8: string (nullable = true)\n\n         \nTo overwrite your schema or change partitioning, please set:\n'.option(\"overwriteSchema\", \"true\")'.\n\nNote that the schema can't be overwritten when using\n'replaceWhere'.\n         \n\n=== Streaming Query ===\nIdentifier: [id = 436fbc25-47c5-41dd-ab69-8ae704819b0f, runId = 3df3de22-f052-40ec-a3aa-e74b240347a3]\nCurrent Committed Offsets: {}\nCurrent Available Offsets: {FileStreamSource[abfss://oea@stoeappexito2023.dfs.core.windows.net/dev/stage1/Transactional/M365/v1.14/Organization/snapshot_batch_data/rundate=2023-04-11 03:03:05]: {\"logOffset\":0}}\n\nCurrent State: ACTIVE\nThread State: RUNNABLE\n\nLogical Plan:\nFileStreamSource[abfss://oea@stoeappexito2023.dfs.core.windows.net/dev/stage1/Transactional/M365/v1.14/Organization/snapshot_batch_data/rundate=2023-04-11 03:03:05]","failureType": "UserError","target": "ingest_insights","details": []}

解决方案

方案1:合并Schema(保留新增的_c8字段)

在Delta写入操作中添加mergeSchema=true参数,允许自动合并新增字段到现有Delta表:

# 修改overwrite方法中的写入代码
df.write.format('delta').mode('overwrite').option('mergeSchema', 'true').save(destination_url)

如果是结构化流写入,在writeStream中添加该参数:

query = streaming_df.writeStream.format('delta').outputMode('append').option('mergeSchema', 'true').trigger(once=True).option('checkpointLocation', checkpoint_path).foreachBatch(wrapped_function).start()

方案2:覆盖原Schema(用K12数据的Schema替换Delta表Schema)

如果确认K12数据的Schema是目标Schema,可使用overwriteSchema=true参数强制覆盖:

df.write.format('delta').mode('overwrite').option('overwriteSchema', 'true').save(destination_url)

注意:使用replaceWhere时无法覆盖Schema,需确保写入逻辑未使用该参数。

方案3:预处理K12数据(移除多余字段)

如果不需要_c8字段,可在写入前过滤掉该字段:

# 仅保留Delta表已有的字段
filtered_df = df.select('_c0', '_c1', '_c2', '_c3', '_c4', '_c5', '_c6', '_c7')
filtered_df.write.format('delta').mode('overwrite').save(destination_url)

额外检查点

  • 确认HED与K12数据集的元数据定义差异,更新管道中的Schema配置以匹配K12数据源
  • 检查metadata.csv中K12对应表的字段定义是否正确

内容的提问来源于stack exchange,提问作者Magoji

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 03:57:00