Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 49 additions & 18 deletions src/io/bam.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
/// BAM output writer with noodles (streaming, unsorted)
use crate::error::Error;
use crate::genome::Genome;
use crate::io::bgzf_writer::BgzfWriter;
use crate::params::Parameters;
use crate::quant::transcriptome::TranscriptomeIndex;
use byteorder::{LittleEndian, WriteBytesExt};
Expand Down Expand Up @@ -45,19 +46,38 @@ fn bgzf_compression(level: i32) -> bgzf::io::writer::CompressionLevel {
}

/// Create a BGZF writer with the given STAR compression level.
fn make_bgzf_writer<W: std::io::Write>(inner: W, compression: i32) -> bgzf::io::Writer<W> {
bgzf::io::writer::Builder::default()
.set_compression_level(bgzf_compression(compression))
.build_from_writer(inner)
///
/// `threads > 1` compresses blocks on that many dedicated worker threads
/// (see [`BgzfWriter`]); the output bytes do not depend on `threads`.
fn make_bgzf_writer<W: Write + Send + 'static>(
inner: W,
compression: i32,
threads: usize,
) -> Result<BgzfWriter<W>, Error> {
Ok(BgzfWriter::new(
inner,
bgzf_compression(compression),
threads,
)?)
}

/// BGZF compression workers for a run: one per `--runThreadN`, like STAR,
/// which deflates its BAM output inside every alignment thread. The workers
/// only use CPU while there are blocks to compress.
fn bgzf_threads(params: &Parameters) -> usize {
params.run_thread_n.get()
}

type BgzfFile = BgzfWriter<BufWriter<File>>;
type BgzfStdout = BgzfWriter<BufWriter<std::io::Stdout>>;

/// BAM file writer (streaming, unsorted)
///
/// This writer streams BAM records directly to disk as they're generated,
/// without buffering or sorting. The output is BGZF-compressed but unsorted.
/// Users can sort the output with `samtools sort` if needed.
pub struct BamWriter {
writer: bam::io::Writer<bgzf::io::Writer<BufWriter<File>>>,
writer: bam::io::Writer<BgzfFile>,
header: sam::Header,
}

Expand All @@ -70,6 +90,7 @@ pub struct SortedBamWriter {
output_path: std::path::PathBuf,
header: sam::Header,
compression: i32,
threads: usize,
limit_bam_sort_ram: u64,
}

Expand All @@ -85,9 +106,10 @@ impl BamWriter {
output_path: &Path,
header: sam::Header,
compression: i32,
threads: usize,
) -> Result<Self, Error> {
let buf_writer = BufWriter::new(File::create(output_path)?);
let mut bgzf = make_bgzf_writer(buf_writer, compression);
let mut bgzf = make_bgzf_writer(buf_writer, compression, threads)?;
write_bam_header_lenient(&mut bgzf, &header, None)?;
let writer = bam::io::Writer::from(bgzf);
Ok(Self { writer, header })
Expand All @@ -99,6 +121,7 @@ impl BamWriter {
output_path,
crate::io::sam::build_sam_header(genome, params)?,
params.out_bam_compression,
bgzf_threads(params),
)
}

Expand All @@ -119,6 +142,7 @@ impl BamWriter {
output_path,
crate::io::sam::build_sam_header_from_refs(refs, params)?,
params.out_bam_compression,
bgzf_threads(params),
)
}

Expand All @@ -135,7 +159,7 @@ impl BamWriter {

/// Flush and close BAM file
pub fn finish(&mut self) -> Result<(), Error> {
self.writer.finish(&self.header)?;
self.writer.get_mut().finish()?;
log::info!("BAM file written successfully");
Ok(())
}
Expand All @@ -154,6 +178,7 @@ impl SortedBamWriter {
output_path: output_path.to_path_buf(),
header,
compression: params.out_bam_compression,
threads: bgzf_threads(params),
limit_bam_sort_ram: params.limit_bam_sort_ram,
})
}
Expand Down Expand Up @@ -198,13 +223,13 @@ impl SortedBamWriter {
});

let buf_writer = BufWriter::new(File::create(&self.output_path)?);
let mut bgzf = make_bgzf_writer(buf_writer, self.compression);
let mut bgzf = make_bgzf_writer(buf_writer, self.compression, self.threads)?;
write_bam_header_lenient(&mut bgzf, &self.header, Some("coordinate"))?;
let mut bam_writer = bam::io::Writer::from(bgzf);
for record in &self.records {
bam_writer.write_alignment_record(&self.header, record)?;
}
bam_writer.finish(&self.header)?;
bam_writer.get_mut().finish()?;
log::info!("Sorted BAM written ({} records)", self.records.len());
Ok(())
}
Expand All @@ -218,15 +243,14 @@ impl SortedBamWriter {
_ => (usize::MAX, 0),
});

let stdout = std::io::stdout();
let buf_writer = BufWriter::new(stdout.lock());
let mut bgzf = make_bgzf_writer(buf_writer, self.compression);
let buf_writer = BufWriter::new(std::io::stdout());
let mut bgzf = make_bgzf_writer(buf_writer, self.compression, self.threads)?;
write_bam_header_lenient(&mut bgzf, &self.header, Some("coordinate"))?;
let mut bam_writer = bam::io::Writer::from(bgzf);
for record in &self.records {
bam_writer.write_alignment_record(&self.header, record)?;
}
bam_writer.finish(&self.header)?;
bam_writer.get_mut().finish()?;
log::info!(
"Sorted BAM written to stdout ({} records)",
self.records.len()
Expand Down Expand Up @@ -371,7 +395,7 @@ fn render_sam_text_lenient(header: &sam::Header, sort_order: Option<&str>) -> Ve

/// Streaming unsorted BAM writer that writes to stdout.
pub struct BamStdoutWriter {
writer: bam::io::Writer<bgzf::io::Writer<BufWriter<std::io::Stdout>>>,
writer: bam::io::Writer<BgzfStdout>,
header: sam::Header,
}

Expand All @@ -381,7 +405,8 @@ impl BamStdoutWriter {
let mut bgzf = make_bgzf_writer(
BufWriter::new(std::io::stdout()),
params.out_bam_compression,
);
bgzf_threads(params),
)?;
write_bam_header_lenient(&mut bgzf, &header, None)?;
let writer = bam::io::Writer::from(bgzf);
Ok(Self { writer, header })
Expand All @@ -395,7 +420,7 @@ impl BamStdoutWriter {
}

pub fn finish(&mut self) -> Result<(), Error> {
self.writer.finish(&self.header)?;
self.writer.get_mut().finish()?;
Ok(())
}
}
Expand All @@ -405,6 +430,7 @@ pub struct SortedBamStdoutWriter {
records: Vec<RecordBuf>,
header: sam::Header,
compression: i32,
threads: usize,
limit_bam_sort_ram: u64,
}

Expand All @@ -415,6 +441,7 @@ impl SortedBamStdoutWriter {
records: Vec::new(),
header,
compression: params.out_bam_compression,
threads: bgzf_threads(params),
limit_bam_sort_ram: params.limit_bam_sort_ram,
})
}
Expand All @@ -441,13 +468,17 @@ impl SortedBamStdoutWriter {
(Some(chr), Some(pos)) => (chr, pos.get()),
_ => (usize::MAX, 0),
});
let mut bgzf = make_bgzf_writer(BufWriter::new(std::io::stdout()), self.compression);
let mut bgzf = make_bgzf_writer(
BufWriter::new(std::io::stdout()),
self.compression,
self.threads,
)?;
write_bam_header_lenient(&mut bgzf, &self.header, Some("coordinate"))?;
let mut bam_writer = bam::io::Writer::from(bgzf);
for record in &self.records {
bam_writer.write_alignment_record(&self.header, record)?;
}
bam_writer.finish(&self.header)?;
bam_writer.get_mut().finish()?;
log::info!(
"Sorted BAM written to stdout ({} records)",
self.records.len()
Expand Down
Loading
Loading