mirror of
https://github.com/Syngnat/GoNavi.git
synced 2026-07-20 04:11:56 +08:00
⚡️ perf(export): 重构大结果集导出链路并支持流式写入
- 新增 ExportFileOptions 统一承载导出格式、进度任务和 XLSX sheet 行数上限 - 查询导出改为流式写入文件,避免一次性缓存整批结果导致高内存占用 - 增加值数组快速路径并复用扫描与写入缓冲,减少逐行 map 分配开销 - 为 ClickHouse、自定义驱动、达梦、SQLServer 和 TDengine 补齐 StreamQuery 支持 - 导出时间字符串仅在形似时间时再解析,避免普通文本被误判改写 - 补充 XLSX 分 sheet、流式导出和基准测试覆盖
This commit is contained in:
@@ -76,6 +76,28 @@ type StatementQueryExecer interface {
|
||||
QueryContext(ctx context.Context, query string) ([]map[string]interface{}, []string, error)
|
||||
}
|
||||
|
||||
// QueryStreamConsumer receives query metadata and rows incrementally.
|
||||
// Implementations can stream rows directly to files to avoid buffering entire result sets in memory.
|
||||
type QueryStreamConsumer interface {
|
||||
SetColumns(columns []string) error
|
||||
ConsumeRow(row map[string]interface{}) error
|
||||
}
|
||||
|
||||
// QueryStreamValueConsumer is an optional fast path for stream consumers that
|
||||
// can consume normalized row values in column order without requiring a
|
||||
// map[string]interface{} allocation per row.
|
||||
type QueryStreamValueConsumer interface {
|
||||
SetColumns(columns []string) error
|
||||
ConsumeRowValues(values []interface{}) error
|
||||
}
|
||||
|
||||
// StreamQueryExecer is an optional interface for drivers or pinned sessions that can
|
||||
// stream query rows incrementally instead of materializing []map rows in memory.
|
||||
type StreamQueryExecer interface {
|
||||
StreamQuery(query string, consumer QueryStreamConsumer) error
|
||||
StreamQueryContext(ctx context.Context, query string, consumer QueryStreamConsumer) error
|
||||
}
|
||||
|
||||
// StatementQueryMessageExecer can run queries on a pinned session and return
|
||||
// extra server messages/notices alongside rows.
|
||||
type StatementQueryMessageExecer interface {
|
||||
@@ -178,6 +200,22 @@ func (e *sqlConnStatementExecer) Query(query string) ([]map[string]interface{},
|
||||
return e.QueryContext(context.Background(), query)
|
||||
}
|
||||
|
||||
func (e *sqlConnStatementExecer) StreamQueryContext(ctx context.Context, query string, consumer QueryStreamConsumer) error {
|
||||
if e == nil || e.conn == nil {
|
||||
return fmt.Errorf("连接未打开")
|
||||
}
|
||||
rows, err := e.conn.QueryContext(ctx, query)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer rows.Close()
|
||||
return streamRowsForDialect(rows, e.scanDialect, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlConnStatementExecer) StreamQuery(query string, consumer QueryStreamConsumer) error {
|
||||
return e.StreamQueryContext(context.Background(), query, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlConnStatementExecer) QueryMultiContext(ctx context.Context, query string) ([]connection.ResultSetData, error) {
|
||||
if e == nil || e.conn == nil {
|
||||
return nil, fmt.Errorf("连接未打开")
|
||||
@@ -275,6 +313,23 @@ func (e *sqlConnTransactionExecer) Query(query string) ([]map[string]interface{}
|
||||
return e.QueryContext(context.Background(), query)
|
||||
}
|
||||
|
||||
func (e *sqlConnTransactionExecer) StreamQueryContext(ctx context.Context, query string, consumer QueryStreamConsumer) error {
|
||||
conn, err := e.activeConn()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rows, err := conn.QueryContext(ctx, query)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer rows.Close()
|
||||
return streamRowsForDialect(rows, e.scanDialect, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlConnTransactionExecer) StreamQuery(query string, consumer QueryStreamConsumer) error {
|
||||
return e.StreamQueryContext(context.Background(), query, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlConnTransactionExecer) QueryMultiContext(ctx context.Context, query string) ([]connection.ResultSetData, error) {
|
||||
conn, err := e.activeConn()
|
||||
if err != nil {
|
||||
@@ -401,6 +456,23 @@ func (e *sqlTxStatementExecer) Query(query string) ([]map[string]interface{}, []
|
||||
return e.QueryContext(context.Background(), query)
|
||||
}
|
||||
|
||||
func (e *sqlTxStatementExecer) StreamQueryContext(ctx context.Context, query string, consumer QueryStreamConsumer) error {
|
||||
tx, err := e.activeTx()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rows, err := tx.QueryContext(ctx, query)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer rows.Close()
|
||||
return streamRows(rows, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlTxStatementExecer) StreamQuery(query string, consumer QueryStreamConsumer) error {
|
||||
return e.StreamQueryContext(context.Background(), query, consumer)
|
||||
}
|
||||
|
||||
func (e *sqlTxStatementExecer) QueryMultiContext(ctx context.Context, query string) ([]connection.ResultSetData, error) {
|
||||
tx, err := e.activeTx()
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user