transaction.go 5.7 KB

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