MCPcopy Create free account
hub / github.com/PostHog/duckgres / batchInsertRows

Method batchInsertRows

server/conn_copy.go:1122–1176  ·  view source on GitHub ↗

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{})

Source from the content-addressed store, hash-verified

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

Callers 2

Calls 5

logQueryStartedMethod · 0.95
logQueryFinishedMethod · 0.95
NowMethod · 0.80
ExecMethod · 0.65
RowsAffectedMethod · 0.65

Tested by

no test coverage detected