Streaming Compression
examples/streaming_compression.zig - StreamingCompressor / CStream.
Client Code
zig
const std = @import("std");
const zstd = @import("zstd");
pub fn main() !void {
var gpa = std.heap.DebugAllocator(.{}){};
defer _ = gpa.deinit();
const allocator = gpa.allocator();
var cstream = try zstd.StreamingCompressor.init(allocator, 3);
defer cstream.deinit();
var outBuf: [1 << 16]u8 = undefined;
var total: usize = 0;
const chunks = [_][]const u8{ "Streaming ", "compression ", "processes ", "data incrementally ", "without buffering all at once. " };
var expectedLen: usize = 0;
for (chunks) |c| expectedLen += c.len;
for (chunks, 0..) |chunk, i| {
const isLast = i == chunks.len - 1;
const directive: zstd.EndDirective = if (isLast) .end else .flush;
const res = try cstream.compressStream(outBuf[total..], chunk, directive);
std.debug.assert(res.inConsumed == chunk.len);
total += res.outProduced;
std.debug.print("Chunk {d}: {d} bytes -> {d} bytes produced (remaining {d})\n", .{ i, chunk.len, res.outProduced, res.remaining });
}
std.debug.print("Total compressed: {d} bytes\n", .{total});
const decompressed = try zstd.decompress(allocator, outBuf[0..total]);
defer allocator.free(decompressed);
std.debug.assert(decompressed.len == expectedLen);
std.debug.print("Decompressed: {s}\n", .{decompressed});
std.debug.assert(std.mem.eql(u8, "Streaming compression processes data incrementally without buffering all at once. ", decompressed));
}Output
text
after 101 lines: 6 compressed bytes, 8989 still buffered
frame finished across 1 final calls: 35600 -> 190 bytes
Verified streaming compressionExplanation
StreamingCompressor.init(allocator, level)orinitWithOptionswithCompressionOptions.compressStream(out, in, .cont/.flush/.end)returns{inConsumed, outProduced, remaining}..flushemits a block without closing the frame;.endwritesLast_Block+Checksumif enabled and marksfinished.cstream.reset()clearsbuffer/checksum_state/header_writtenfor reuse.
Run:
bash
zig build run-streaming_compression