transaction.go 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273
  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. tableRowItems := make([]*request.InsertTableRowItem, 0)
  60. for _, reqTableItem := range req.Items {
  61. tableRows := make([]*request.TableRow, 0)
  62. for _, reqTableRow := range reqTableItem.TableRows {
  63. if reqTableRow == nil {
  64. continue
  65. }
  66. tableRow, err := reqTableRow.ToDPSTableRow()
  67. if err != nil {
  68. return "", err
  69. }
  70. tableRows = append(tableRows, tableRow)
  71. }
  72. tableRowItems = append(tableRowItems, &request.InsertTableRowItem{
  73. TablePrefixWithSchema: reqTableItem.TablePrefixWithSchema,
  74. Version: reqTableItem.Version,
  75. KeyColumns: reqTableItem.KeyColumns,
  76. TableRows: tableRows,
  77. })
  78. }
  79. err := tx.stream.Send(&request.TransactionOperation{
  80. Request: &request.TransactionOperation_InsertBatchTxRequest{
  81. InsertBatchTxRequest: &request.InsertBatchTxRequest{
  82. Items: tableRowItems,
  83. UserID: req.UserID,
  84. },
  85. }})
  86. if err != nil {
  87. return "", err
  88. }
  89. txResponse, err := tx.stream.Recv()
  90. if err != nil {
  91. return "", err
  92. }
  93. return txResponse.Statement, nil
  94. }
  95. func (tx *Transaction) DeleteTx(req *client.DeleteRequest) (string, error) {
  96. tx.client.transactionMutex.Lock()
  97. defer tx.client.transactionMutex.Unlock()
  98. if tx.client.conn == nil {
  99. return "", nil
  100. }
  101. err := tx.stream.Send(&request.TransactionOperation{
  102. Request: &request.TransactionOperation_DeleteTxRequest{
  103. DeleteTxRequest: &request.DeleteTxRequest{
  104. TablePrefixWithSchema: req.TablePrefixWithSchema,
  105. Version: req.Version,
  106. KeyColumns: req.KeyValues.Columns(),
  107. KeyValues: req.KeyValues.Values(),
  108. UserID: req.UserID,
  109. },
  110. }})
  111. if err != nil {
  112. return "", err
  113. }
  114. txResponse, err := tx.stream.Recv()
  115. if err != nil {
  116. return "", err
  117. }
  118. return txResponse.Statement, nil
  119. }
  120. func (tx *Transaction) DeleteWhereTx(req *client.DeleteWhereRequest) (string, error) {
  121. if req.Where == nil {
  122. return "", errors.New("没有传递Where条件")
  123. }
  124. tx.client.transactionMutex.Lock()
  125. defer tx.client.transactionMutex.Unlock()
  126. if tx.client.conn == nil {
  127. return "", nil
  128. }
  129. whereJsonBytes, err := req.Where.ToJson()
  130. if err != nil {
  131. return "", err
  132. }
  133. err = tx.stream.Send(&request.TransactionOperation{
  134. Request: &request.TransactionOperation_DeleteWhereTxRequest{
  135. DeleteWhereTxRequest: &request.DeleteWhereTxRequest{
  136. TablePrefixWithSchema: req.TablePrefixWithSchema,
  137. Version: req.Version,
  138. KeyColumns: req.KeyColumns,
  139. Where: whereJsonBytes,
  140. UserID: req.UserID,
  141. },
  142. }})
  143. if err != nil {
  144. return "", err
  145. }
  146. txResponse, err := tx.stream.Recv()
  147. if err != nil {
  148. return "", err
  149. }
  150. return txResponse.Statement, nil
  151. }
  152. func (tx *Transaction) UpdateTx(req *client.UpdateRequest) (string, error) {
  153. if req.NewTableRow == nil {
  154. return "", nil
  155. }
  156. tx.client.transactionMutex.Lock()
  157. defer tx.client.transactionMutex.Unlock()
  158. if tx.client.conn == nil {
  159. return "", nil
  160. }
  161. reqNewTableRow, err := req.NewTableRow.ToDPSTableRow()
  162. if err != nil {
  163. return "", err
  164. }
  165. err = tx.stream.Send(&request.TransactionOperation{
  166. Request: &request.TransactionOperation_UpdateTxRequest{
  167. UpdateTxRequest: &request.UpdateTxRequest{
  168. TablePrefixWithSchema: req.TablePrefixWithSchema,
  169. Version: req.Version,
  170. KeyColumns: req.KeyValues.Columns(),
  171. KeyValues: req.KeyValues.Values(),
  172. NewTableRow: reqNewTableRow,
  173. UserID: req.UserID,
  174. },
  175. }})
  176. if err != nil {
  177. return "", err
  178. }
  179. txResponse, err := tx.stream.Recv()
  180. if err != nil {
  181. return "", err
  182. }
  183. return txResponse.Statement, nil
  184. }
  185. func (tx *Transaction) UpdateWhereTx(req *client.UpdateWhereRequest) (string, error) {
  186. if req.Where == nil {
  187. return "", errors.New("没有传递Where条件")
  188. }
  189. if req.NewTableRow == nil {
  190. return "", nil
  191. }
  192. tx.client.transactionMutex.Lock()
  193. defer tx.client.transactionMutex.Unlock()
  194. if tx.client.conn == nil {
  195. return "", nil
  196. }
  197. whereJsonBytes, err := req.Where.ToJson()
  198. if err != nil {
  199. return "", err
  200. }
  201. reqNewTableRow, err := req.NewTableRow.ToDPSTableRow()
  202. if err != nil {
  203. return "", err
  204. }
  205. err = tx.stream.Send(&request.TransactionOperation{
  206. Request: &request.TransactionOperation_UpdateWhereTxRequest{
  207. UpdateWhereTxRequest: &request.UpdateWhereTxRequest{
  208. TablePrefixWithSchema: req.TablePrefixWithSchema,
  209. Version: req.Version,
  210. KeyColumns: req.KeyColumns,
  211. Where: whereJsonBytes,
  212. NewTableRow: reqNewTableRow,
  213. UserID: req.UserID,
  214. },
  215. }})
  216. if err != nil {
  217. return "", err
  218. }
  219. txResponse, err := tx.stream.Recv()
  220. if err != nil {
  221. return "", err
  222. }
  223. return txResponse.Statement, nil
  224. }