蛋蛋星球RabbitMq消费项目
25개 이상의 토픽을 선택하실 수 없습니다. Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

150 lines
4.1 KiB

  1. package db
  2. import (
  3. "database/sql"
  4. "fmt"
  5. "os"
  6. "time"
  7. _ "github.com/go-sql-driver/mysql" //必须导入mysql驱动,否则会panic
  8. "xorm.io/xorm"
  9. "xorm.io/xorm/log"
  10. "applet/app/cfg"
  11. "applet/app/utils/logx"
  12. )
  13. var Db *xorm.Engine
  14. // 根据DB配置文件初始化数据库
  15. func InitDB(c *cfg.DBCfg) error {
  16. var (
  17. err error
  18. f *os.File
  19. )
  20. //创建Orm引擎
  21. if Db, err = xorm.NewEngine("mysql", fmt.Sprintf("%s:%s@tcp(%s)/%s?charset=utf8mb4", c.User, c.Psw, c.Host, c.Name)); err != nil {
  22. return err
  23. }
  24. Db.SetConnMaxLifetime(c.MaxLifetime * time.Second) //设置最长连接时间
  25. Db.SetMaxOpenConns(c.MaxOpenConns) //设置最大打开连接数
  26. Db.SetMaxIdleConns(c.MaxIdleConns) //设置连接池的空闲数大小
  27. if err = Db.Ping(); err != nil { //尝试ping数据库
  28. return err
  29. }
  30. if c.ShowLog { //根据配置文件设置日志
  31. Db.ShowSQL(true) //设置是否打印sql
  32. Db.Logger().SetLevel(0) //设置日志等级
  33. //修改日志文件存放路径文件名是%s.log
  34. path := fmt.Sprintf(c.Path, c.Name)
  35. f, err = os.OpenFile(path, os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0777)
  36. if err != nil {
  37. os.RemoveAll(c.Path)
  38. if f, err = os.OpenFile(c.Path, os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0777); err != nil {
  39. return err
  40. }
  41. }
  42. logger := log.NewSimpleLogger(f)
  43. logger.ShowSQL(true)
  44. Db.SetLogger(logger)
  45. }
  46. return nil
  47. }
  48. var DbIm *xorm.Engine
  49. // 根据DB配置文件初始化数据库
  50. func InitImDB(c *cfg.DBCfg) error {
  51. var (
  52. err error
  53. f *os.File
  54. )
  55. //创建Orm引擎
  56. if DbIm, err = xorm.NewEngine("mysql", fmt.Sprintf("%s:%s@tcp(%s)/%s?charset=utf8mb4", c.User, c.Psw, c.Host, c.Name)); err != nil {
  57. return err
  58. }
  59. DbIm.SetConnMaxLifetime(c.MaxLifetime * time.Second) //设置最长连接时间
  60. DbIm.SetMaxOpenConns(c.MaxOpenConns) //设置最大打开连接数
  61. DbIm.SetMaxIdleConns(c.MaxIdleConns) //设置连接池的空闲数大小
  62. if err = DbIm.Ping(); err != nil { //尝试ping数据库
  63. return err
  64. }
  65. if c.ShowLog { //根据配置文件设置日志
  66. DbIm.ShowSQL(true) //设置是否打印sql
  67. DbIm.Logger().SetLevel(0) //设置日志等级
  68. //修改日志文件存放路径文件名是%s.log
  69. path := fmt.Sprintf(c.Path, c.Name)
  70. f, err = os.OpenFile(path, os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0777)
  71. if err != nil {
  72. os.RemoveAll(c.Path)
  73. if f, err = os.OpenFile(c.Path, os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0777); err != nil {
  74. return err
  75. }
  76. }
  77. logger := log.NewSimpleLogger(f)
  78. logger.ShowSQL(true)
  79. DbIm.SetLogger(logger)
  80. }
  81. return nil
  82. }
  83. /********************************************* 公用方法 *********************************************/
  84. // 数据批量插入
  85. func DbInsertBatch(Db *xorm.Engine, m ...interface{}) error {
  86. if len(m) == 0 {
  87. return nil
  88. }
  89. id, err := Db.Insert(m...)
  90. if id == 0 || err != nil {
  91. return logx.Warn("cannot insert data :", err)
  92. }
  93. return nil
  94. }
  95. // QueryNativeString 查询原生sql
  96. func QueryNativeString(Db *xorm.Engine, sql string, args ...interface{}) ([]map[string]string, error) {
  97. results, err := Db.SQL(sql, args...).QueryString()
  98. return results, err
  99. }
  100. // UpdateComm common update
  101. func UpdateComm(Db *xorm.Engine, id interface{}, model interface{}) (int64, error) {
  102. row, err := Db.ID(id).Update(model)
  103. return row, err
  104. }
  105. // InsertComm common insert
  106. func InsertComm(Db *xorm.Engine, model interface{}) (int64, error) {
  107. row, err := Db.InsertOne(model)
  108. return row, err
  109. }
  110. // ExecuteOriginalSql 执行原生sql
  111. func ExecuteOriginalSql(session *xorm.Session, sql string) (sql.Result, error) {
  112. result, err := session.Exec(sql)
  113. if err != nil {
  114. _ = logx.Warn(err)
  115. return nil, err
  116. }
  117. return result, nil
  118. }
  119. // GetComm
  120. // payload *model
  121. // return *model,has,err
  122. func GetComm(Db *xorm.Engine, model interface{}) (interface{}, bool, error) {
  123. has, err := Db.Get(model)
  124. if err != nil {
  125. _ = logx.Warn(err)
  126. return nil, false, err
  127. }
  128. return model, has, nil
  129. }
  130. // InsertCommWithSession common insert
  131. func InsertCommWithSession(session *xorm.Session, model interface{}) (int64, error) {
  132. row, err := session.InsertOne(model)
  133. return row, err
  134. }