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 filesconst streams = fileHelper.readAsStream({ root: 'assets', glob: '**/*' });for (const stream of streams) {// Access raw Node.js streamsconst { reader, writer } = stream;// Process chunk by chunk using readerlet content = '';for await (const chunk of reader) {content += chunk.toString();}// Transform and write using writerconst processed = transformContent(content);writer.write(processed);writer.end();}return { directory: input.writeDir };});
filename="processor.py"
from cyanprintsdk import start_processor_with_fnfrom cyanprintsdk.domain.processor.input import ProcessorInputfrom cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelperasync def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):# Get streams for large filesstreams = fileHelper.read_as_stream(root='assets', glob='**/*')for stream in streams:# Access raw Python streamsreader = stream.readerwriter = stream.writer# Process chunk by chunk using readercontent = ''async for chunk in reader:content += chunk.decode('utf-8')# Transform and write using writerprocessed = 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 filesvar streams = fileHelper.ReadAsStream(root: "assets", glob: "**/*");foreach (var stream in streams){// Access raw .NET streamsvar reader = stream.Reader;var writer = stream.Writer;// Process chunk by chunk using readerusing var streamReader = new StreamReader(reader);var content = await streamReader.ReadToEndAsync();// Transform and write using writervar 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 linefor 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_fnfrom cyanprintsdk.domain.processor.input import ProcessorInputfrom cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelperasync def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):streams = fileHelper.read_as_stream(root='data', glob='**/*.csv')for stream in streams:reader = stream.readerwriter = stream.writerlines = []# Process line by lineasync 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 lineusing 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 processingconst chunks: Buffer[] = [];for await (const chunk of reader) {chunks.push(chunk);}// Combine and processconst buffer = Buffer.concat(chunks);const processed = transformImage(buffer);writer.write(processed.toString('base64'));writer.end();}return { directory: input.writeDir };});
filename="processor.py"
import base64from cyanprintsdk import start_processor_with_fnfrom cyanprintsdk.domain.processor.input import ProcessorInputfrom cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelperasync def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):streams = fileHelper.read_as_stream(root='images', glob='**/*')for stream in streams:reader = stream.readerwriter = stream.writer# Collect all chunks for binary processingchunks = []async for chunk in reader:chunks.append(chunk)# Combine and processbuffer = 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 processingusing var memoryStream = new MemoryStream();await reader.CopyToAsync(memoryStream);var buffer = memoryStream.ToArray();// Combine and processvar 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 arrivefor 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 refrom cyanprintsdk import start_processor_with_fnfrom cyanprintsdk.domain.processor.input import ProcessorInputfrom cyanprintsdk.domain.core.fs.cyan_fs_helper import CyanFileHelperasync def my_processor(input: ProcessorInput, fileHelper: CyanFileHelper):streams = fileHelper.read_as_stream(root='logs', glob='*.log')for stream in streams:reader = stream.readerwriter = stream.writerresult = ''# Process chunks as they arriveasync for chunk in reader:transformed = chunk.decode('utf-8')transformed = re.sub(r'ERROR', 'ERR', transformed)transformed = re.sub(r'WARNING', 'WARN', transformed)result += transformedwriter.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 arriveusing 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.