【问题标题】:Databricks Delta table Alter column for decimal(10,0) to decimal(38,18) conversion not workingDatabricks Delta 表更改十进制(10,0)到十进制(38,18)转换的列不起作用
【发布时间】:2021-12-04 14:22:23
【问题描述】:

在 Databricks 中,表是使用模式 json 定义创建的。

用于创建表的schema json

{
  "fields": [
    {
      "metadata": {},
      "name": "username",
      "nullable": true,
      "type": "string"
    },
    {
      "metadata": {},
      "name": "department",
      "nullable": true,
      "type": "string"
    },
    {
      "metadata": {},
      "name": "income",
      "nullable": true,
      "type": "decimal(38,18)"
    }
    ],
    "type" :"struct"
}

以下代码创建表格

  ...
   # the schema json file is placed in a location and using it to create table
   with open('/dbfs/FileStore/my-schema/{0}.json'.format(tbl_name), 'r') as f:
          #  data files
          tbl_schema = T.StructType.fromJson(json.loads(f.read()))
          tbl_df = spark.createDataFrame([], tbl_schema)
          tbl_df.write.format("delta").save(tbl_path)
      # create table.
      spark.sql("CREATE TABLE {0} USING DELTA LOCATION '{1}'".format(tbl_name, tbl_path))
   ...

当我 describe 表时,我看到 DecimalType(10,0) 表示收入字段。

我正在使用 readstreams 从 ORC 文件中读取数据,其中使用了 Decimal(38,18),能够在数据帧中 printScehma()。

我正在使用 Spark 结构化流写入流,它使用 UPSERT(合并到),使用 foreachBatch()。

类似于此链接中的python示例https://docs.azuredatabricks.net/_static/notebooks/merge-in-streaming.html

问题:

每当我运行表时,都不会插入数据。也没有调试日志消息。

我想问题可能是由于表中的字段收入为 DecimalType(10,0),而数据框为 DecimalType(38,18)。

所以我正在尝试更改字段,但无法做到。我正在使用以下命令。

%sql ALTER TABLE mytable ALTER income TYPE decimal(38,18);
com.databricks.backend.common.rpc.DatabricksExceptions$SQLExecutionException: org.apache.spark.sql.AnalysisException: ALTER TABLE CHANGE COLUMN is not supported for changing column 'income' with type 'DecimalType(10,0) (nullable = true)' to 'income' with type 'DecimalType(38,18) (nullable = true)'
    at com.databricks.sql.transaction.tahoe.DeltaErrors$.alterTableChangeColumnException(DeltaErrors.scala:478)
    at com.databricks.sql.transaction.tahoe.commands.AlterTableChangeColumnDeltaCommand.verifyColumnChange(alterDeltaTableCommands.scala:412)
    at com.databricks.sql.transaction.tahoe.commands.AlterTableChangeColumnDeltaCommand.$anonfun$run$15(alterDeltaTableCommands.scala:295)
    at com.databricks.sql.transaction.tahoe.schema.SchemaUtils$.transform$1(SchemaUtils.scala:762)
    at com.databricks.sql.transaction.tahoe.schema.SchemaUtils$.transformColumnsStructs(SchemaUtils.scala:781)
    at com.databricks.sql.transaction.tahoe.commands.AlterTableChangeColumnDeltaCommand.$anonfun$run$14(alterDeltaTableCommands.scala:292)
    at com.databricks.spark.util.FrameProfiler$.record(FrameProfiler.scala:80)
    at com.databricks.sql.transaction.tahoe.metering.DeltaLogging.$anonfun$recordDeltaOperation$5(DeltaLogging.scala:122

如果我删除 TYPE ,我会得到以下异常

com.databricks.backend.common.rpc.DatabricksExceptions$SQLExecutionException: org.apache.spark.sql.catalyst.parser.ParseException: 
mismatched input 'decimal' expecting {<EOF>, ';'}(line 1, pos 46)

== SQL ==
ALTER TABLE ahm_db.message_raw ALTER income decimal(38,18)
----------------------------------------------^^^

    at org.apache.spark.sql.catalyst.parser.ParseException.withCommand(ParseDriver.scala:265)
    at org.apache.spark.sql.catalyst.parser.AbstractSqlParser.parse(ParseDriver.scala:134)
    at org.apache.spark.sql.execution.SparkSqlParser.parse(SparkSqlParser.scala:64)
    at org.apache.spark.sql.catalyst.parser.AbstractSqlParser.parsePlan(ParseDriver.scala:85)
    at com.databricks.sql.parser.DatabricksSqlParser.$anonfun$parsePlan$1(DatabricksSqlParser.scala:67)
    at com.databricks.sql.parser.DatabricksSqlParser.parse(DatabricksSqlParser.scala:87)

【问题讨论】:

  • 作为一种解决方法,我使用了dataframe.writestream.option('mergeSchema','true')。同样在我的情况下,我使用的是dataframe.writestream.option('mergeSchema','true').option('checkpointLocation',"/path/to/_checkpoint/')..。因为我从 blob 存储中执行了 ORC 文件,并清理了一些文件。检查点在内存中跟踪它,每次我执行笔记本时,它都会尝试执行上一个操作。

标签: spark-streaming databricks azure-databricks delta-lake


【解决方案1】:

因为在我的案例中是 delta Lake 商店,所以启用了 options('checkpoint','/_checkpoint') 选项。数据参考仍然可用。

在开发过程中,删除表并优化使用 vaccum 之后。

%sql delete from <my-table-name>
spark.conf.set('spark.databricks.delta.retentionDurationCheck.enabled','false')
# use spark.conf.get('property') to check default or current value
# Per documentation setting retention duration less than 7 days is not a recommended practice, but depends on requirements
%sql
VACUUM <my-table-name> RETAIN 0 HOURS
%sql
drop table <my-table-name>
  • 从 ADLS Gen2 存储帐户容器中删除了 _checkpoint 文件夹。

完成上述步骤后,重新运行 json 特定架构并应用它。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-01-09
    • 2016-01-17
    • 1970-01-01
    • 2016-05-18
    • 2017-05-05
    • 2018-02-23
    相关资源
    最近更新 更多