LogoCyanPrint
ProcessorsHow-To Guides

Stream Large Files

Stream large files for memory efficiency

Stream Large Files

Use readAsStream() to process large files chunk by chunk without loading the entire file into memory.

When to Use

Best for:

  • Large binary files
  • Files larger than available memory
  • Memory-constrained environments
  • Processing files line by line

Usage

filename="processor.ts"
StartProcessorWithLambda(async (input, fileHelper) => {
// Get streams for large files
const streams = fileHelper.readAsStream({ root: 'assets', glob: '**/*' });
for (const stream of streams) {
// Access raw Node.js streams
const { reader, writer } = stream;
// Process chunk by chunk using reader
let content = '';
for await (const chunk of reader) {
content += chunk.toString();
}
// Transform and write using writer
const processed = transformContent(content);
writer.write(processed);
writer.end();
}
return { directory: input.writeDir };
});
filename="processor.py"
from cyanprintsdk import start_processor_with_fn
from cyanprintsdk.domain.processor.input import ProcessorInput
from cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelper
async def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):
# Get streams for large files
streams = fileHelper.read_as_stream(root='assets', glob='**/*')
for stream in streams:
# Access raw Python streams
reader = stream.reader
writer = stream.writer
# Process chunk by chunk using reader
content = ''
async for chunk in reader:
content += chunk.decode('utf-8')
# Transform and write using writer
processed = transform_content(content)
writer.write(processed.encode('utf-8'))
await writer.drain()
writer.close()
return {'directory': input.write_directory}
start_processor_with_fn(my_processor)
filename="Processor.cs"
using System.IO;
using sulfone_helium;
using sulfone_helium.Domain.Core.FileSystem;
using sulfone_helium.Domain.Processor;
CyanEngine.StartProcessor(args, async (input, fileHelper) =>
{
// Get streams for large files
var streams = fileHelper.ReadAsStream(root: "assets", glob: "**/*");
foreach (var stream in streams)
{
// Access raw .NET streams
var reader = stream.Reader;
var writer = stream.Writer;
// Process chunk by chunk using reader
using var streamReader = new StreamReader(reader);
var content = await streamReader.ReadToEndAsync();
// Transform and write using writer
var processed = TransformContent(content);
using var streamWriter = new StreamWriter(writer);
await streamWriter.WriteAsync(processed);
}
return new ProcessorOutput { Directory = input.WriteDirectory };
});

VirtualFileStream Properties

Prop

Type

VirtualFileStream exposes raw Node.js streams. Use standard Node.js stream APIs to read from reader and write to writer.

Example: Line-by-Line Processing

filename="processor.ts"
StartProcessorWithLambda(async (input, fileHelper) => {
const streams = fileHelper.readAsStream({ root: 'data', glob: '**/*.csv' });
for (const stream of streams) {
const { reader, writer } = stream;
const lines: string[] = [];
// Process line by line
for await (const chunk of reader) {
const text = chunk.toString();
const newLines = text.split('\n');
for (const line of newLines) {
if (line.trim()) {
lines.push(processLine(line));
}
}
}
writer.write(lines.join('\n'));
writer.end();
}
return { directory: input.writeDir };
});
filename="processor.py"
from cyanprintsdk import start_processor_with_fn
from cyanprintsdk.domain.processor.input import ProcessorInput
from cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelper
async def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):
streams = fileHelper.read_as_stream(root='data', glob='**/*.csv')
for stream in streams:
reader = stream.reader
writer = stream.writer
lines = []
# Process line by line
async for chunk in reader:
text = chunk.decode('utf-8')
new_lines = text.split('\n')
for line in new_lines:
if line.strip():
lines.append(process_line(line))
writer.write('\n'.join(lines).encode('utf-8'))
await writer.drain()
writer.close()
return {'directory': input.write_directory}
start_processor_with_fn(my_processor)
filename="Processor.cs"
using System.IO;
using sulfone_helium;
using sulfone_helium.Domain.Core.FileSystem;
using sulfone_helium.Domain.Processor;
CyanEngine.StartProcessor(args, async (input, fileHelper) =>
{
var streams = fileHelper.ReadAsStream(root: "data", glob: "**/*.csv");
foreach (var stream in streams)
{
var reader = stream.Reader;
var writer = stream.Writer;
var lines = new List<string>();
// Process line by line
using var streamReader = new StreamReader(reader);
string? line;
while ((line = await streamReader.ReadLineAsync()) != null)
{
if (!string.IsNullOrWhiteSpace(line))
{
lines.Add(ProcessLine(line));
}
}
using var streamWriter = new StreamWriter(writer);
await streamWriter.WriteAsync(string.Join('\n', lines));
}
return new ProcessorOutput { Directory = input.WriteDirectory };
});

Example: Binary File Processing

filename="processor.ts"
StartProcessorWithLambda(async (input, fileHelper) => {
const streams = fileHelper.readAsStream({ root: 'images', glob: '**/*' });
for (const stream of streams) {
const { reader, writer } = stream;
// Collect all chunks for binary processing
const chunks: Buffer[] = [];
for await (const chunk of reader) {
chunks.push(chunk);
}
// Combine and process
const buffer = Buffer.concat(chunks);
const processed = transformImage(buffer);
writer.write(processed.toString('base64'));
writer.end();
}
return { directory: input.writeDir };
});
filename="processor.py"
import base64
from cyanprintsdk import start_processor_with_fn
from cyanprintsdk.domain.processor.input import ProcessorInput
from cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelper
async def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):
streams = fileHelper.read_as_stream(root='images', glob='**/*')
for stream in streams:
reader = stream.reader
writer = stream.writer
# Collect all chunks for binary processing
chunks = []
async for chunk in reader:
chunks.append(chunk)
# Combine and process
buffer = b''.join(chunks)
processed = transform_image(buffer)
writer.write(base64.b64encode(processed))
await writer.drain()
writer.close()
return {'directory': input.write_directory}
start_processor_with_fn(my_processor)
filename="Processor.cs"
using System.IO;
using sulfone_helium;
using sulfone_helium.Domain.Core.FileSystem;
using sulfone_helium.Domain.Processor;
CyanEngine.StartProcessor(args, async (input, fileHelper) =>
{
var streams = fileHelper.ReadAsStream(root: "images", glob: "**/*");
foreach (var stream in streams)
{
var reader = stream.Reader;
var writer = stream.Writer;
// Collect all chunks for binary processing
using var memoryStream = new MemoryStream();
await reader.CopyToAsync(memoryStream);
var buffer = memoryStream.ToArray();
// Combine and process
var processed = TransformImage(buffer);
using var streamWriter = new StreamWriter(writer);
await streamWriter.WriteAsync(Convert.ToBase64String(processed));
}
return new ProcessorOutput { Directory = input.WriteDirectory };
});

Example: Transform While Streaming

filename="processor.ts"
import { Transform } from 'stream';
StartProcessorWithLambda(async (input, fileHelper) => {
const streams = fileHelper.readAsStream({ root: 'logs', glob: '*.log' });
for (const stream of streams) {
const { reader, writer } = stream;
let result = '';
// Process chunks as they arrive
for await (const chunk of reader) {
const transformed = chunk
.toString()
.replace(/ERROR/g, 'ERR')
.replace(/WARNING/g, 'WARN');
result += transformed;
}
writer.write(result);
writer.end();
}
return { directory: input.writeDir };
});
filename="processor.py"
import re
from cyanprintsdk import start_processor_with_fn
from cyanprintsdk.domain.processor.input import ProcessorInput
from cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelper
async def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):
streams = fileHelper.read_as_stream(root='logs', glob='*.log')
for stream in streams:
reader = stream.reader
writer = stream.writer
result = ''
# Process chunks as they arrive
async for chunk in reader:
transformed = chunk.decode('utf-8')
transformed = re.sub(r'ERROR', 'ERR', transformed)
transformed = re.sub(r'WARNING', 'WARN', transformed)
result += transformed
writer.write(result.encode('utf-8'))
await writer.drain()
writer.close()
return {'directory': input.write_directory}
start_processor_with_fn(my_processor)
filename="Processor.cs"
using System.IO;
using System.Text.RegularExpressions;
using sulfone_helium;
using sulfone_helium.Domain.Core.FileSystem;
using sulfone_helium.Domain.Processor;
CyanEngine.StartProcessor(args, async (input, fileHelper) =>
{
var streams = fileHelper.ReadAsStream(root: "logs", glob: "*.log");
foreach (var stream in streams)
{
var reader = stream.Reader;
var writer = stream.Writer;
var result = new StringBuilder();
// Process chunks as they arrive
using var streamReader = new StreamReader(reader);
var content = await streamReader.ReadToEndAsync();
var transformed = Regex.Replace(content, "ERROR", "ERR");
transformed = Regex.Replace(transformed, "WARNING", "WARN");
using var streamWriter = new StreamWriter(writer);
await streamWriter.WriteAsync(transformed);
}
return new ProcessorOutput { Directory = input.WriteDirectory };
});

Streaming requires accumulating content before writing. For truly massive files that don't fit in memory, consider using fileHelper.copy() to pass them through unchanged.