appendWithDuckDBAppender uses the DuckDB Appender API for fast bulk inserts. Only works for full-column inserts (no column subset). batchInsertRows inserts rows using batched multi-row INSERT statements. Used as fallback when Appender can't be used (column subsets, unsupported types).
(tableName, columnList string, cols []string, rows [][]interface{})
| 1120 | break // EOF or error — stop reading |
| 1121 | } |
| 1122 | if len(record) != len(cols) { |
| 1123 | c.logger().Warn("COPY FROM STDIN (BLOB fallback) skipping row with wrong field count.", "expected", len(cols), "got", len(record)) |
| 1124 | continue |
| 1125 | } |
| 1126 | row := make([]interface{}, len(cols)) |
| 1127 | for j, field := range record { |
| 1128 | if field == opts.NullString { |
| 1129 | row[j] = nil |
| 1130 | } else if isBlobCol[j] { |
| 1131 | row[j] = []byte(field) |
| 1132 | } else { |
| 1133 | row[j] = field |
| 1134 | } |
| 1135 | } |
| 1136 | rows = append(rows, row) |
| 1137 | } |
| 1138 | |
| 1139 | if len(rows) == 0 { |
| 1140 | _ = c.writeCommandComplete("COPY 0") |
| 1141 | c.logQuery(copyStartTime, query, query, "COPY", 0, 0, "", "", "simple") |
| 1142 | _ = c.writeReadyForQuery(c.txStatus) |
| 1143 | _ = c.flushWriter() |
| 1144 | return nil |
| 1145 | } |
| 1146 | |
| 1147 | loadStart := time.Now() |
| 1148 | rowCount, err := c.batchInsertRows(opts.TableName, opts.ColumnList, cols, rows) |
| 1149 | if err != nil { |
| 1150 | c.logger().Error("COPY FROM STDIN (BLOB fallback) INSERT failed.", "error", err) |
| 1151 | errMsg := fmt.Sprintf("COPY failed: %v", err) |
| 1152 | c.sendError("ERROR", "22P02", errMsg) |
| 1153 | c.setTxError() |
| 1154 | c.logQuery(copyStartTime, query, query, "COPY", 0, int64(rowCount), "22P02", errMsg, "simple") |
| 1155 | _ = c.writeReadyForQuery(c.txStatus) |
| 1156 | _ = c.flushWriter() |
| 1157 | return nil |
| 1158 | } |
| 1159 | |
| 1160 | totalElapsed := time.Since(copyStartTime) |
| 1161 | loadElapsed := time.Since(loadStart) |
| 1162 | c.logger().Info("COPY FROM STDIN (BLOB fallback) completed.", "rows", rowCount, "total_duration", totalElapsed, "load_duration", loadElapsed) |
| 1163 | |
| 1164 | _ = c.writeCommandComplete(fmt.Sprintf("COPY %d", rowCount)) |
| 1165 | c.logQuery(copyStartTime, query, query, "COPY", 0, int64(rowCount), "", "", "simple") |
| 1166 | _ = c.writeReadyForQuery(c.txStatus) |
| 1167 | _ = c.flushWriter() |
| 1168 | return nil |
| 1169 | |
| 1170 | case wire.MsgCopyFail: |
| 1171 | errMsg := string(bytes.TrimRight(body, "\x00")) |
| 1172 | exception := fmt.Sprintf("COPY failed: %s", errMsg) |
| 1173 | c.sendError("ERROR", "57014", exception) |
| 1174 | c.setTxError() |
| 1175 | c.logQuery(copyStartTime, query, query, "COPY", 0, 0, "57014", exception, "simple") |
| 1176 | _ = c.writeReadyForQuery(c.txStatus) |
| 1177 | _ = c.flushWriter() |
| 1178 | return nil |
| 1179 |
no test coverage detected