( ctx *sql.Context, conn *stdsql.Conn, tx *stdsql.Tx, table tableIdentifier, appender *DeltaAppender, stats *FlushStats, )
| 315 | } |
| 316 | |
| 317 | func (c *DeltaController) handleDeleteOnly( |
| 318 | ctx *sql.Context, |
| 319 | conn *stdsql.Conn, |
| 320 | tx *stdsql.Tx, |
| 321 | table tableIdentifier, |
| 322 | appender *DeltaAppender, |
| 323 | stats *FlushStats, |
| 324 | ) error { |
| 325 | // Ignore all but the primary key fields |
| 326 | viewName, release, err := c.prepareArrowView(ctx, conn, table, appender, 0, getPrimaryKeyIndices(appender)) |
| 327 | if err != nil { |
| 328 | return err |
| 329 | } |
| 330 | defer release() |
| 331 | |
| 332 | qualifiedTableName := catalog.ConnectIdentifiersANSI(table.dbName, table.tableName) |
| 333 | pk := getPrimaryKeyStruct(appender.BaseSchema()) |
| 334 | |
| 335 | // Perform direct DELETE without deduplication |
| 336 | deleteSQL := "DELETE FROM " + qualifiedTableName + |
| 337 | " WHERE " + pk + " IN (SELECT " + pk + " FROM " + viewName + ")" |
| 338 | result, err := tx.ExecContext(ctx, deleteSQL) |
| 339 | if err != nil { |
| 340 | return err |
| 341 | } |
| 342 | |
| 343 | affected, err := result.RowsAffected() |
| 344 | if err != nil { |
| 345 | return err |
| 346 | } |
| 347 | stats.Deletions += affected |
| 348 | stats.DeltaSize += affected |
| 349 | |
| 350 | if log := ctx.GetLogger(); log.Logger.IsLevelEnabled(logrus.DebugLevel) { |
| 351 | log.WithFields(logrus.Fields{ |
| 352 | "db": table.dbName, |
| 353 | "table": table.tableName, |
| 354 | "rows": affected, |
| 355 | }).Debug("Deleted") |
| 356 | } |
| 357 | |
| 358 | return nil |
| 359 | } |
| 360 | |
| 361 | func (c *DeltaController) handleZeroDelete( |
| 362 | ctx *sql.Context, |
no test coverage detected