(fn string, comp compressionType)
| 253 | } |
| 254 | |
| 255 | func (s *CompressedInput) parseFileTyped(fn string, comp compressionType) { |
| 256 | |
| 257 | ctx := log.WithFields(log.Fields{"f": "compressedInput.parseFile", "fn": fn}) |
| 258 | stream, sz, lastModified, url, err := s.Opener(fn) |
| 259 | |
| 260 | stream = s.stats.NewStatsReader(stream, sz) |
| 261 | if err != nil { |
| 262 | log.WithFields(log.Fields{"f": "compressedInput.parseFile", "fn": fn}).WithError(err).Error("Error while opening stream") |
| 263 | return |
| 264 | } |
| 265 | defer stream.Close() |
| 266 | |
| 267 | var r io.Reader |
| 268 | |
| 269 | switch comp { |
| 270 | case gzipCompression: |
| 271 | if sz > 1000000 { |
| 272 | rgz, err := newFastGzReader(stream) |
| 273 | if err != nil { |
| 274 | // Sometimes the fast gz reader fails to initialize due to |
| 275 | // memory pressure. We'd still like to run so try the |
| 276 | // slower (and less memory hungry) gzip. |
| 277 | ctx.WithError(err).Error("error initializing fast gzip, will attempt slow gzip") |
| 278 | r, err = gzip.NewReader(stream) |
| 279 | if err != nil { |
| 280 | ctx.WithError(err).Fatal("both fast and slow gzip readers failed to initialize") |
| 281 | return |
| 282 | } |
| 283 | } else { |
| 284 | defer rgz.Close() |
| 285 | r = rgz |
| 286 | } |
| 287 | } else { |
| 288 | rgz, err := gzip.NewReader(stream) |
| 289 | if err != nil { |
| 290 | ctx.WithError(err).Fatal("error initializing gzip") |
| 291 | return |
| 292 | } |
| 293 | defer rgz.Close() |
| 294 | r = rgz |
| 295 | } |
| 296 | case zstdCompression: |
| 297 | rzst := zstd.NewReader(stream) |
| 298 | defer rzst.Release() |
| 299 | r = rzst |
| 300 | default: |
| 301 | ctx.WithError(err).Fatal("Unknown compression type specified.") |
| 302 | } |
| 303 | |
| 304 | ctx.Info("begin reading") |
| 305 | |
| 306 | rbuf := bufio.NewReaderSize(r, kChunkBuffer) |
| 307 | |
| 308 | for atomic.LoadInt64(&s.stopping) == 0 { |
| 309 | bakerData := s.pool.Get().(*baker.Data) |
| 310 | bakerData.Meta = baker.Metadata{ |
| 311 | MetadataLastModified: lastModified, |
| 312 | MetadataURL: url, |
no test coverage detected