transaction.go 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321
  1. package dpsv1
  2. import (
  3. "errors"
  4. "git.sxidc.com/service-supports/dps-sdk/client"
  5. "git.sxidc.com/service-supports/dps-sdk/pb/v1"
  6. "git.sxidc.com/service-supports/dps-sdk/pb/v1/request"
  7. )
  8. type Transaction struct {
  9. stream v1.CommandService_TransactionClient
  10. client *Client
  11. }
  12. func (tx *Transaction) InsertTx(req *client.InsertRequest) (string, error) {
  13. if req.TableRow == nil {
  14. return "", nil
  15. }
  16. tx.client.transactionMutex.Lock()
  17. defer tx.client.transactionMutex.Unlock()
  18. if tx.client.conn == nil {
  19. return "", nil
  20. }
  21. var err error
  22. defer func() {
  23. if err != nil {
  24. innerErr := tx.stream.CloseSend()
  25. if innerErr != nil {
  26. panic(innerErr)
  27. }
  28. }
  29. }()
  30. reqTableRow, err := req.TableRow.ToDPSTableRow()
  31. if err != nil {
  32. return "", err
  33. }
  34. err = tx.stream.Send(&request.TransactionOperation{
  35. Request: &request.TransactionOperation_InsertTxRequest{
  36. InsertTxRequest: &request.InsertTxRequest{
  37. TablePrefixWithSchema: req.TablePrefixWithSchema,
  38. Version: req.Version,
  39. KeyColumns: req.KeyColumns,
  40. TableRow: reqTableRow,
  41. UserID: req.UserID,
  42. },
  43. }})
  44. if err != nil {
  45. return "", err
  46. }
  47. txResponse, err := tx.stream.Recv()
  48. if err != nil {
  49. return "", err
  50. }
  51. return txResponse.Statement, nil
  52. }
  53. func (tx *Transaction) InsertBatchTx(req *client.InsertBatchRequest) (string, error) {
  54. tx.client.transactionMutex.Lock()
  55. defer tx.client.transactionMutex.Unlock()
  56. if tx.client.conn == nil {
  57. return "", nil
  58. }
  59. var err error
  60. defer func() {
  61. if err != nil {
  62. innerErr := tx.stream.CloseSend()
  63. if innerErr != nil {
  64. panic(innerErr)
  65. }
  66. }
  67. }()
  68. tableRowItems := make([]*request.InsertTableRowItem, 0)
  69. for _, reqTableItem := range req.Items {
  70. tableRows := make([]*request.TableRow, 0)
  71. for _, reqTableRow := range reqTableItem.TableRows {
  72. if reqTableRow == nil {
  73. continue
  74. }
  75. tableRow, err := reqTableRow.ToDPSTableRow()
  76. if err != nil {
  77. return "", err
  78. }
  79. tableRows = append(tableRows, tableRow)
  80. }
  81. tableRowItems = append(tableRowItems, &request.InsertTableRowItem{
  82. TablePrefixWithSchema: reqTableItem.TablePrefixWithSchema,
  83. Version: reqTableItem.Version,
  84. KeyColumns: reqTableItem.KeyColumns,
  85. TableRows: tableRows,
  86. })
  87. }
  88. err = tx.stream.Send(&request.TransactionOperation{
  89. Request: &request.TransactionOperation_InsertBatchTxRequest{
  90. InsertBatchTxRequest: &request.InsertBatchTxRequest{
  91. Items: tableRowItems,
  92. UserID: req.UserID,
  93. },
  94. }})
  95. if err != nil {
  96. return "", err
  97. }
  98. txResponse, err := tx.stream.Recv()
  99. if err != nil {
  100. return "", err
  101. }
  102. return txResponse.Statement, nil
  103. }
  104. func (tx *Transaction) DeleteTx(req *client.DeleteRequest) (string, error) {
  105. tx.client.transactionMutex.Lock()
  106. defer tx.client.transactionMutex.Unlock()
  107. if tx.client.conn == nil {
  108. return "", nil
  109. }
  110. var err error
  111. defer func() {
  112. if err != nil {
  113. innerErr := tx.stream.CloseSend()
  114. if innerErr != nil {
  115. panic(innerErr)
  116. }
  117. }
  118. }()
  119. err = tx.stream.Send(&request.TransactionOperation{
  120. Request: &request.TransactionOperation_DeleteTxRequest{
  121. DeleteTxRequest: &request.DeleteTxRequest{
  122. TablePrefixWithSchema: req.TablePrefixWithSchema,
  123. Version: req.Version,
  124. KeyValues: req.KeyValues,
  125. UserID: req.UserID,
  126. },
  127. }})
  128. if err != nil {
  129. return "", err
  130. }
  131. txResponse, err := tx.stream.Recv()
  132. if err != nil {
  133. return "", err
  134. }
  135. return txResponse.Statement, nil
  136. }
  137. func (tx *Transaction) DeleteWhereTx(req *client.DeleteWhereRequest) (string, error) {
  138. tx.client.transactionMutex.Lock()
  139. defer tx.client.transactionMutex.Unlock()
  140. if tx.client.conn == nil {
  141. return "", nil
  142. }
  143. if req.Where == nil {
  144. return "", errors.New("没有传递Where条件")
  145. }
  146. var err error
  147. defer func() {
  148. if err != nil {
  149. innerErr := tx.stream.CloseSend()
  150. if innerErr != nil {
  151. panic(innerErr)
  152. }
  153. }
  154. }()
  155. whereJsonBytes, err := req.Where.ToJson()
  156. if err != nil {
  157. return "", err
  158. }
  159. err = tx.stream.Send(&request.TransactionOperation{
  160. Request: &request.TransactionOperation_DeleteWhereTxRequest{
  161. DeleteWhereTxRequest: &request.DeleteWhereTxRequest{
  162. TablePrefixWithSchema: req.TablePrefixWithSchema,
  163. Version: req.Version,
  164. KeyColumns: req.KeyColumns,
  165. Where: whereJsonBytes,
  166. UserID: req.UserID,
  167. },
  168. }})
  169. if err != nil {
  170. return "", err
  171. }
  172. txResponse, err := tx.stream.Recv()
  173. if err != nil {
  174. return "", err
  175. }
  176. return txResponse.Statement, nil
  177. }
  178. func (tx *Transaction) DeleteBatchTx(req *client.DeleteBatchRequest) (string, error) {
  179. tx.client.transactionMutex.Lock()
  180. defer tx.client.transactionMutex.Unlock()
  181. if tx.client.conn == nil {
  182. return "", nil
  183. }
  184. var err error
  185. defer func() {
  186. if err != nil {
  187. innerErr := tx.stream.CloseSend()
  188. if innerErr != nil {
  189. panic(innerErr)
  190. }
  191. }
  192. }()
  193. tableRowItems := make([]*request.DeleteTableRowItem, 0)
  194. for _, reqTableItem := range req.Items {
  195. items := make([]*request.DeleteItem, 0)
  196. for _, reqKeyValues := range reqTableItem.KeyValues {
  197. items = append(items, &request.DeleteItem{
  198. KeyValues: reqKeyValues,
  199. })
  200. }
  201. tableRowItems = append(tableRowItems, &request.DeleteTableRowItem{
  202. TablePrefixWithSchema: reqTableItem.TablePrefixWithSchema,
  203. Version: reqTableItem.Version,
  204. Items: items,
  205. })
  206. }
  207. err = tx.stream.Send(&request.TransactionOperation{
  208. Request: &request.TransactionOperation_DeleteBatchTxRequest{
  209. DeleteBatchTxRequest: &request.DeleteBatchTxRequest{
  210. Items: tableRowItems,
  211. UserID: req.UserID,
  212. },
  213. }})
  214. if err != nil {
  215. return "", err
  216. }
  217. txResponse, err := tx.stream.Recv()
  218. if err != nil {
  219. return "", err
  220. }
  221. return txResponse.Statement, nil
  222. }
  223. func (tx *Transaction) UpdateTx(req *client.UpdateRequest) (string, error) {
  224. if req.NewTableRow == nil {
  225. return "", nil
  226. }
  227. tx.client.transactionMutex.Lock()
  228. defer tx.client.transactionMutex.Unlock()
  229. if tx.client.conn == nil {
  230. return "", nil
  231. }
  232. var err error
  233. defer func() {
  234. if err != nil {
  235. innerErr := tx.stream.CloseSend()
  236. if innerErr != nil {
  237. panic(innerErr)
  238. }
  239. }
  240. }()
  241. reqNewTableRow, err := req.NewTableRow.ToDPSTableRow()
  242. if err != nil {
  243. return "", err
  244. }
  245. err = tx.stream.Send(&request.TransactionOperation{
  246. Request: &request.TransactionOperation_UpdateTxRequest{
  247. UpdateTxRequest: &request.UpdateTxRequest{
  248. TablePrefixWithSchema: req.TablePrefixWithSchema,
  249. Version: req.Version,
  250. KeyValues: req.KeyValues,
  251. NewTableRow: reqNewTableRow,
  252. UserID: req.UserID,
  253. },
  254. }})
  255. if err != nil {
  256. return "", err
  257. }
  258. txResponse, err := tx.stream.Recv()
  259. if err != nil {
  260. return "", err
  261. }
  262. return txResponse.Statement, nil
  263. }