u
This commit is contained in:
1 parent
f373cf956d
commit
91ec5e268c
12 files changed
+114
-64
No files matched your search
File renamed without changes.
@@ -1,7 +1,7 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
error2 "base-framework/pkg/utils/error"
|
||||
error2 "base-framework/pkg/error"
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"io"
|
||||
|
||||
@@ -2,37 +2,67 @@ package db
|
||||
|
||||
import (
|
||||
"base-framework/pkg/config"
|
||||
"base-framework/pkg/reponse"
|
||||
"base-framework/pkg/router"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"fmt"
|
||||
)
|
||||
|
||||
// Client 封装数据库连接对象
|
||||
type Client struct {
|
||||
Conn *sql.DB
|
||||
conn *sql.DB
|
||||
tx *sql.Tx
|
||||
}
|
||||
|
||||
// NewClient 根据 context 获取 orgId 并返回 Client 对象
|
||||
func NewClient(c *router.Context) *Client {
|
||||
// JdbcTemplate 创建 JdbcTemplate
|
||||
func JdbcTemplate(c *router.Context) (*Client, error) {
|
||||
val, ok := c.Get("orgId")
|
||||
fmt.Println(val)
|
||||
if !ok {
|
||||
reponse.Error(c).Code(reponse.CodeInvalidOrgCode).Send()
|
||||
return nil
|
||||
return nil, errors.New("missing orgId in context")
|
||||
}
|
||||
|
||||
orgId, ok := val.(string)
|
||||
if !ok || orgId == "" {
|
||||
reponse.Error(c).Code(reponse.CodeInvalidOrgCode).Send()
|
||||
return nil
|
||||
return nil, errors.New("invalid orgId in context")
|
||||
}
|
||||
|
||||
conn, ok := config.GetDB(orgId)
|
||||
if !ok || conn == nil {
|
||||
reponse.Error(c).Code(reponse.CodeInvalidOrgCode).Send()
|
||||
return nil
|
||||
return nil, fmt.Errorf("no db connection found for orgId=%s", orgId)
|
||||
}
|
||||
return &Client{conn: conn}, nil
|
||||
}
|
||||
|
||||
// ------------------------ 内部方法 ------------------------
|
||||
|
||||
// 执行查询,返回 *sql.Rows
|
||||
func (c *Client) query(query string, args ...any) (*sql.Rows, error) {
|
||||
if c.tx != nil {
|
||||
return c.tx.Query(query, args...)
|
||||
}
|
||||
return c.conn.Query(query, args...)
|
||||
}
|
||||
|
||||
// 执行执行类语句(insert/update/delete)
|
||||
func (c *Client) exec(query string, args ...any) (sql.Result, error) {
|
||||
if c.tx != nil {
|
||||
return c.tx.Exec(query, args...)
|
||||
}
|
||||
return c.conn.Exec(query, args...)
|
||||
}
|
||||
|
||||
// WithTransaction 自动处理事务提交或回滚
|
||||
func (c *Client) WithTransaction(fn func(txClient *Client) error) error {
|
||||
if c.tx != nil {
|
||||
// 已经在事务中,直接执行
|
||||
return fn(c)
|
||||
}
|
||||
|
||||
return &Client{Conn: conn}
|
||||
tx, err := c.conn.Begin()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
txClient := &Client{conn: c.conn, tx: tx}
|
||||
if err := fn(txClient); err != nil {
|
||||
_ = tx.Rollback()
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
@@ -1 +1,5 @@
|
||||
package db
|
||||
|
||||
import "database/sql"
|
||||
|
||||
func (c *Client) Delete(query string, args ...any) (sql.Result, error) { return c.exec(query, args...) }
|
||||
@@ -1 +1,15 @@
|
||||
package db
|
||||
|
||||
import "database/sql"
|
||||
|
||||
func (c *Client) Insert(query string, args ...any) (sql.Result, error) { return c.exec(query, args...) }
|
||||
|
||||
// BatchInsert 批量插入,传入多组参数
|
||||
func (c *Client) BatchInsert(query string, params [][]any) error {
|
||||
for _, args := range params {
|
||||
if _, err := c.Insert(query, args...); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1,62 +1,63 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
)
|
||||
import "database/sql"
|
||||
|
||||
func (c *Client) QueryRows(query string) ([]map[string]interface{}, error) {
|
||||
if c.Conn == nil {
|
||||
return nil, fmt.Errorf("database connection is nil")
|
||||
}
|
||||
|
||||
rows, err := c.Conn.Query(query)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("query error: %v", err)
|
||||
}
|
||||
defer func(rows *sql.Rows) {
|
||||
err := rows.Close()
|
||||
if err != nil {
|
||||
fmt.Println("关闭资源失败")
|
||||
}
|
||||
}(rows)
|
||||
|
||||
// 获取列名
|
||||
cols, err := rows.Columns()
|
||||
func (c *Client) Select(query string, args ...any) ([]map[string]any, error) {
|
||||
rows, err := c.query(query, args...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer func(rows *sql.Rows) {
|
||||
_ = rows.Close()
|
||||
}(rows)
|
||||
|
||||
var result []map[string]interface{}
|
||||
cols, _ := rows.Columns()
|
||||
var result []map[string]any
|
||||
|
||||
for rows.Next() {
|
||||
// 创建扫描用的切片
|
||||
columnPointers := make([]interface{}, len(cols))
|
||||
columnValues := make([]interface{}, len(cols))
|
||||
for i := range columnPointers {
|
||||
columnPointers[i] = &columnValues[i]
|
||||
columns := make([]any, len(cols))
|
||||
columnPointers := make([]any, len(cols))
|
||||
for i := range columns {
|
||||
columnPointers[i] = &columns[i]
|
||||
}
|
||||
|
||||
if err := rows.Scan(columnPointers...); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
rowMap := make(map[string]interface{})
|
||||
rowMap := make(map[string]any)
|
||||
for i, colName := range cols {
|
||||
val := columnValues[i]
|
||||
if b, ok := val.([]byte); ok {
|
||||
rowMap[colName] = string(b)
|
||||
} else {
|
||||
rowMap[colName] = val
|
||||
}
|
||||
rowMap[colName] = columns[i]
|
||||
}
|
||||
|
||||
result = append(result, rowMap)
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
if err := rows.Err(); err != nil {
|
||||
func (c *Client) SelectOne(query string, args ...any) (map[string]any, error) {
|
||||
rows, err := c.query(query, args...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer func(rows *sql.Rows) {
|
||||
_ = rows.Close()
|
||||
}(rows)
|
||||
|
||||
if !rows.Next() {
|
||||
return nil, sql.ErrNoRows
|
||||
}
|
||||
|
||||
cols, _ := rows.Columns()
|
||||
columns := make([]any, len(cols))
|
||||
columnPointers := make([]any, len(cols))
|
||||
for i := range columns {
|
||||
columnPointers[i] = &columns[i]
|
||||
}
|
||||
if err := rows.Scan(columnPointers...); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return result, nil
|
||||
rowMap := make(map[string]any)
|
||||
for i, colName := range cols {
|
||||
rowMap[colName] = columns[i]
|
||||
}
|
||||
return rowMap, nil
|
||||
}
|
||||
@@ -1 +1,5 @@
|
||||
package db
|
||||
|
||||
import "database/sql"
|
||||
|
||||
func (c *Client) Update(query string, args ...any) (sql.Result, error) { return c.exec(query, args...) }
|
||||
Reference in new issue
Block a user