compress.go 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194
  1. package sarama
  2. import (
  3. "bytes"
  4. "compress/gzip"
  5. "fmt"
  6. "sync"
  7. snappy "github.com/eapache/go-xerial-snappy"
  8. "github.com/pierrec/lz4"
  9. )
  10. var (
  11. lz4WriterPool = sync.Pool{
  12. New: func() interface{} {
  13. return lz4.NewWriter(nil)
  14. },
  15. }
  16. gzipWriterPool = sync.Pool{
  17. New: func() interface{} {
  18. return gzip.NewWriter(nil)
  19. },
  20. }
  21. gzipWriterPoolForCompressionLevel1 = sync.Pool{
  22. New: func() interface{} {
  23. gz, err := gzip.NewWriterLevel(nil, 1)
  24. if err != nil {
  25. panic(err)
  26. }
  27. return gz
  28. },
  29. }
  30. gzipWriterPoolForCompressionLevel2 = sync.Pool{
  31. New: func() interface{} {
  32. gz, err := gzip.NewWriterLevel(nil, 2)
  33. if err != nil {
  34. panic(err)
  35. }
  36. return gz
  37. },
  38. }
  39. gzipWriterPoolForCompressionLevel3 = sync.Pool{
  40. New: func() interface{} {
  41. gz, err := gzip.NewWriterLevel(nil, 3)
  42. if err != nil {
  43. panic(err)
  44. }
  45. return gz
  46. },
  47. }
  48. gzipWriterPoolForCompressionLevel4 = sync.Pool{
  49. New: func() interface{} {
  50. gz, err := gzip.NewWriterLevel(nil, 4)
  51. if err != nil {
  52. panic(err)
  53. }
  54. return gz
  55. },
  56. }
  57. gzipWriterPoolForCompressionLevel5 = sync.Pool{
  58. New: func() interface{} {
  59. gz, err := gzip.NewWriterLevel(nil, 5)
  60. if err != nil {
  61. panic(err)
  62. }
  63. return gz
  64. },
  65. }
  66. gzipWriterPoolForCompressionLevel6 = sync.Pool{
  67. New: func() interface{} {
  68. gz, err := gzip.NewWriterLevel(nil, 6)
  69. if err != nil {
  70. panic(err)
  71. }
  72. return gz
  73. },
  74. }
  75. gzipWriterPoolForCompressionLevel7 = sync.Pool{
  76. New: func() interface{} {
  77. gz, err := gzip.NewWriterLevel(nil, 7)
  78. if err != nil {
  79. panic(err)
  80. }
  81. return gz
  82. },
  83. }
  84. gzipWriterPoolForCompressionLevel8 = sync.Pool{
  85. New: func() interface{} {
  86. gz, err := gzip.NewWriterLevel(nil, 8)
  87. if err != nil {
  88. panic(err)
  89. }
  90. return gz
  91. },
  92. }
  93. gzipWriterPoolForCompressionLevel9 = sync.Pool{
  94. New: func() interface{} {
  95. gz, err := gzip.NewWriterLevel(nil, 9)
  96. if err != nil {
  97. panic(err)
  98. }
  99. return gz
  100. },
  101. }
  102. )
  103. func compress(cc CompressionCodec, level int, data []byte) ([]byte, error) {
  104. switch cc {
  105. case CompressionNone:
  106. return data, nil
  107. case CompressionGZIP:
  108. var (
  109. err error
  110. buf bytes.Buffer
  111. writer *gzip.Writer
  112. )
  113. switch level {
  114. case CompressionLevelDefault:
  115. writer = gzipWriterPool.Get().(*gzip.Writer)
  116. defer gzipWriterPool.Put(writer)
  117. writer.Reset(&buf)
  118. case 1:
  119. writer = gzipWriterPoolForCompressionLevel1.Get().(*gzip.Writer)
  120. defer gzipWriterPoolForCompressionLevel1.Put(writer)
  121. writer.Reset(&buf)
  122. case 2:
  123. writer = gzipWriterPoolForCompressionLevel2.Get().(*gzip.Writer)
  124. defer gzipWriterPoolForCompressionLevel2.Put(writer)
  125. writer.Reset(&buf)
  126. case 3:
  127. writer = gzipWriterPoolForCompressionLevel3.Get().(*gzip.Writer)
  128. defer gzipWriterPoolForCompressionLevel3.Put(writer)
  129. writer.Reset(&buf)
  130. case 4:
  131. writer = gzipWriterPoolForCompressionLevel4.Get().(*gzip.Writer)
  132. defer gzipWriterPoolForCompressionLevel4.Put(writer)
  133. writer.Reset(&buf)
  134. case 5:
  135. writer = gzipWriterPoolForCompressionLevel5.Get().(*gzip.Writer)
  136. defer gzipWriterPoolForCompressionLevel5.Put(writer)
  137. writer.Reset(&buf)
  138. case 6:
  139. writer = gzipWriterPoolForCompressionLevel6.Get().(*gzip.Writer)
  140. defer gzipWriterPoolForCompressionLevel6.Put(writer)
  141. writer.Reset(&buf)
  142. case 7:
  143. writer = gzipWriterPoolForCompressionLevel7.Get().(*gzip.Writer)
  144. defer gzipWriterPoolForCompressionLevel7.Put(writer)
  145. writer.Reset(&buf)
  146. case 8:
  147. writer = gzipWriterPoolForCompressionLevel8.Get().(*gzip.Writer)
  148. defer gzipWriterPoolForCompressionLevel8.Put(writer)
  149. writer.Reset(&buf)
  150. case 9:
  151. writer = gzipWriterPoolForCompressionLevel9.Get().(*gzip.Writer)
  152. defer gzipWriterPoolForCompressionLevel9.Put(writer)
  153. writer.Reset(&buf)
  154. default:
  155. writer, err = gzip.NewWriterLevel(&buf, level)
  156. if err != nil {
  157. return nil, err
  158. }
  159. }
  160. if _, err := writer.Write(data); err != nil {
  161. return nil, err
  162. }
  163. if err := writer.Close(); err != nil {
  164. return nil, err
  165. }
  166. return buf.Bytes(), nil
  167. case CompressionSnappy:
  168. return snappy.Encode(data), nil
  169. case CompressionLZ4:
  170. writer := lz4WriterPool.Get().(*lz4.Writer)
  171. defer lz4WriterPool.Put(writer)
  172. var buf bytes.Buffer
  173. writer.Reset(&buf)
  174. if _, err := writer.Write(data); err != nil {
  175. return nil, err
  176. }
  177. if err := writer.Close(); err != nil {
  178. return nil, err
  179. }
  180. return buf.Bytes(), nil
  181. case CompressionZSTD:
  182. return zstdCompress(nil, data)
  183. default:
  184. return nil, PacketEncodingError{fmt.Sprintf("unsupported compression codec (%d)", cc)}
  185. }
  186. }