Skip to content

Native Parquet writes on Spark 3.4/3.5 leave INSERT INTO targets stale until REFRESH TABLE #6315

Description

@andygrove

Describe the bug

#3521 was closed by #5763, which fixed it on Spark 4.0+ by leaving InsertIntoHadoopFsRelationCommand in the plan. On Spark 3.4 and 3.5, CometNativeWriteExec still replaces the whole DataWritingCommandExec, so none of the steps InsertIntoHadoopFsRelationCommand.run performs after the write happen: fileIndex.refresh(), cacheManager.recacheByPath, the partition refresh and CommandUtils.updateTableStats. A table written by a native INSERT INTO ... SELECT reads back empty until something refreshes it.

Steps to reproduce

On main at 7649361 with -Pspark-3.5. The same happens with #4746 applied.

val nativeWrite = Seq(
  CometConf.COMET_NATIVE_PARQUET_WRITE_ENABLED.key -> "true",
  CometConf.COMET_OPERATOR_DATA_WRITING_COMMAND_ALLOW_INCOMPAT.key -> "true",
  CometConf.COMET_EXEC_ENABLED.key -> "true")

sql("CREATE TABLE src(id bigint, name string) USING parquet")
sql("CREATE TABLE tgt(id bigint, name string) USING parquet")
withSQLConf(CometConf.COMET_ENABLED.key -> "false") {
  sql("INSERT INTO src VALUES (1, 'a'), (2, 'b')")
}
withSQLConf(nativeWrite: _*) {
  sql("INSERT INTO tgt SELECT id, name FROM src") // plans as CometNativeWrite
}
spark.table("tgt").count() // 0
spark.catalog.refreshTable("tgt")
spark.table("tgt").count() // 2

Expected behavior

spark.table("tgt").count() returns 2 straight after the insert, as it does with Spark's own writer and with Comet on Spark 4.0+.

Additional context

CometParquetWriterSuite gates "INSERT INTO ... SELECT is visible to subsequent reads" and "INSERT INTO ... SELECT writes the target table's column names" to Spark 4.0+ because of this. Both should run on 3.x once it's fixed.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:writerNative Parquet writerbugSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions