Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
2cf26ce
Add shared-filesystem parallel linclust foundations.
bbuschkaemper Jul 29, 2026
53efff4
Add distributed k-mer shuffle.
bbuschkaemper Jul 29, 2026
ec7c139
Add distributed reduce.
bbuschkaemper Jul 29, 2026
67fce30
Add distributed alignment.
bbuschkaemper Jul 30, 2026
21a07be
Add distributed align2clust.
bbuschkaemper Jul 30, 2026
fa62a16
Add k-mer extraction waves (to manage within scratch budget).
bbuschkaemper Jul 30, 2026
0842bb4
Add clusterstsv conversion.
bbuschkaemper Jul 30, 2026
1962152
Seed the pass-2 filter gate with the pair's diagonal.
bbuschkaemper Jul 30, 2026
9cbe59c
Various smaller bugfixes.
bbuschkaemper Jul 31, 2026
2c31978
Fix handling of duplicate edges across blocks.
bbuschkaemper Aug 3, 2026
bc53556
Fix invalid inline documentation.
bbuschkaemper Aug 3, 2026
3742f78
Changed parallel linclust param from mpi-runner to runner.
bbuschkaemper Aug 3, 2026
01d70bf
Fix candidate edges header bug.
bbuschkaemper Aug 3, 2026
8cb0a4c
Update in-line comments.
bbuschkaemper Aug 3, 2026
ac83fe3
Fix name spelling.
bbuschkaemper Aug 3, 2026
80aded2
Fix cicd for greedycluster and kmer partition test.
bbuschkaemper Aug 3, 2026
3f1a0aa
Fix resume crashes. Fix smaller bugs.
bbuschkaemper Aug 3, 2026
e957f2d
Revert unnecessary fsync fix (out of scope, requires too much resourc…
bbuschkaemper Aug 3, 2026
0b78481
Fix silent error in merge cluster parallel. Fix memory budget over-al…
bbuschkaemper Aug 3, 2026
f5b53de
More fixes.
bbuschkaemper Aug 4, 2026
8fd9702
More fixes.
bbuschkaemper Aug 4, 2026
db1f30e
Add varint codec and length-rank table
bbuschkaemper Aug 5, 2026
6f0cfcd
Pack the k-mer and candidate-edge bucket formats
bbuschkaemper Aug 5, 2026
e612926
Fix quadratic growth when decoding packed blocks
bbuschkaemper Aug 5, 2026
648e438
Group oversized partitions in k-mer slices
bbuschkaemper Aug 5, 2026
a014722
Wake the work-queue heartbeat on a condition variable
bbuschkaemper Aug 5, 2026
76463df
Add packed-edge codec tests and align-bucket skew reporting
bbuschkaemper Aug 5, 2026
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
1 change: 1 addition & 0 deletions data/workflow/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ set(GENERATED_WORKFLOWS
workflow/map.sh
workflow/rbh.sh
workflow/linclust.sh
workflow/linclustparallel.sh
workflow/clustering.sh
workflow/cascaded_clustering.sh
workflow/update_clustering.sh
Expand Down
363 changes: 363 additions & 0 deletions data/workflow/linclustparallel.sh

Large diffs are not rendered by default.

10 changes: 10 additions & 0 deletions src/CommandDeclarations.h
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,15 @@ extern int convertkb(int argc, const char **argv, const Command& command);
extern int convertmsa(int argc, const char **argv, const Command& command);
extern int convertprofiledb(int argc, const char **argv, const Command& command);
extern int createdb(int argc, const char **argv, const Command& command);
extern int createdbparallel(int argc, const char **argv, const Command& command);
extern int kmermatcherparallel(int argc, const char **argv, const Command& command);
extern int kmerreduceparallel(int argc, const char **argv, const Command& command);
extern int alignparallel(int argc, const char **argv, const Command& command);
extern int greedycluster(int argc, const char **argv, const Command& command);
extern int mergeclusterparallel(int argc, const char **argv, const Command& command);
extern int createrepdb(int argc, const char **argv, const Command& command);
extern int translatecluster(int argc, const char **argv, const Command& command);
extern int translatekeys(int argc, const char **argv, const Command& command);
extern int makepaddedseqdb(int argc, const char **argv, const Command& command);
extern int createindex(int argc, const char **argv, const Command& command);
extern int createlinindex(int argc, const char **argv, const Command& command);
Expand Down Expand Up @@ -75,6 +84,7 @@ extern int lca(int argc, const char **argv, const Command& command);
extern int lcaalign(int argc, const char **argv, const Command& command);
extern int taxonomyreport(int argc, const char **argv, const Command& command);
extern int linclust(int argc, const char **argv, const Command& command);
extern int linclustparallel(int argc, const char **argv, const Command& command);
extern int map(int argc, const char **argv, const Command& command);
extern int renamedbkeys(int argc, const char **argv, const Command& command);
extern int majoritylca(int argc, const char **argv, const Command& command);
Expand Down
132 changes: 132 additions & 0 deletions src/MMseqsBase.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,18 @@ std::vector<Command> baseCommands = {
"<i:fastaFile1[.gz|.bz2]> ... <i:fastaFileN[.gz|.bz2]>|<i:stdin> <o:sequenceDB>",
CITATION_MMSEQS2, {{"fast[a|q]File[.gz|bz2]|stdin", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA | DbType::VARIADIC, &DbValidator::flatfileStdinAndGeneric },
{"sequenceDB", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"createdbparallel", createdbparallel, &par.createdbparallel, COMMAND_DATABASE_CREATION,
"Convert FASTA file(s) to a length-ranked sequence DB with many workers",
"# Build a sequence DB with sequences ordered longest first and dense keys.\n"
"# Every worker runs this identical command line and coordinates through\n"
"# <sequenceDB>.coord; workers may join late, die and be restarted.\n"
"mmseqs createdbparallel seq.fasta sequenceDB\n\n"
"# Run it from several nodes against the same shared filesystem\n"
"srun -N 8 mmseqs createdbparallel seq.fasta sequenceDB --chunk-size 1G\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:fastaFile1> ... <i:fastaFileN> <o:sequenceDB>",
CITATION_MMSEQS2, {{"fastaFile", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA | DbType::VARIADIC, &DbValidator::flatfile },
{"sequenceDB", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"makepaddedseqdb", makepaddedseqdb, &par.makepaddedseqdb, COMMAND_HIDDEN,
"Generate a padded sequence DB",
"Generate a padded sequence DB",
Expand Down Expand Up @@ -297,6 +309,32 @@ std::vector<Command> baseCommands = {
CITATION_MMSEQS2|CITATION_LINCLUST, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"clusterDB", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::clusterDb },
{"tmpDir", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::directory }}},
{"linclustparallel", linclustparallel, &par.linclustparallelworkflow, COMMAND_MAIN,
"Linclust across many nodes over a shared filesystem",
"# The same clustering as linclust *for the coverage modes it supports*,\n"
"# computed by many independent worker processes that coordinate only\n"
"# through files: no MPI, no rank argument, and no node-to-node\n"
"# communication. Every worker of a stage runs the identical command line,\n"
"# so a stage maps onto a Slurm array job and workers may join late, die,\n"
"# or be restarted. Re-running resumes.\n\n"
"# Supported subset: --cov-mode 1 (target, the default here) or 2 (query).\n"
"# Stock linclust defaults to --cov-mode 0, which selects SET_COVER\n"
"# clustering and the count-table grouping rounds; neither is implemented\n"
"# here, so mode 0 is refused rather than silently answered differently.\n"
"# Nucleotide input is refused for the same reason (stock routes it to the\n"
"# v1 path, which has no counterpart here).\n\n"
"# On one node\n"
"mmseqs linclustparallel sequenceDB clusters.tsv tmp\n\n"
"# Across 64 nodes, bounding scratch at 100 TB\n"
"mmseqs linclustparallel sequenceDB clusters.tsv tmp \\\n"
" --runner \"srun -n 64\" --scratch-budget 100T --split-memory-limit 700G\n\n"
"# Output is representative<TAB>member in accessions, not a cluster DB:\n"
"# a per-key index is state no single node can hold at this scale.\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:fastaFile|sequenceDB> <o:clusterTsv> <tmpDir>",
CITATION_MMSEQS2|CITATION_LINCLUST, {{"fastaFile|sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfileAndSequenceDb },
{"clusterTsv", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"tmpDir", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::directory }}},
{"cluster", clusteringworkflow, &par.clusterworkflow, COMMAND_MAIN,
"Slower, sensitive clustering",
"# Cascaded clustering of FASTA file\n"
Expand Down Expand Up @@ -667,6 +705,100 @@ std::vector<Command> baseCommands = {
"<i:sequenceDB> <o:prefilterDB>",
CITATION_MMSEQS2,{{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"prefilterDB", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::prefilterDb }}},
{"kmermatcherparallel", kmermatcherparallel, &par.kmermatcherparallel, COMMAND_PREFILTER,
"Shuffle a sequence DB's k-mers into partition buckets with many workers",
"# Scan the sequence DB once and write every k-mer into the bucket of its\n"
"# partition. Every worker runs this identical command line and coordinates\n"
"# through <kmerDir>/coord; workers may join late, die and be restarted.\n"
"mmseqs kmermatcherparallel sequenceDB kmerDir\n\n"
"# Run it from several nodes against the same shared filesystem\n"
"srun -N 8 mmseqs kmermatcherparallel sequenceDB kmerDir --scratch-budget 100T\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <o:kmerDir>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"kmerDir", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::directory }}},
{"kmerreduceparallel", kmerreduceparallel, &par.kmerreduceparallel, COMMAND_PREFILTER,
"Group shuffled k-mer partitions into candidate edges with many workers",
"# Read back the partitions kmermatcherparallel wrote, group each one, and\n"
"# write its (representative, member) candidate edges as packed binary.\n"
"# Workers claim whole partitions from <edgeDir>/coord.\n"
"mmseqs kmerreduceparallel sequenceDB kmerDir edgeDir\n\n"
"# Run it from several nodes against the same shared filesystem\n"
"srun -N 8 mmseqs kmerreduceparallel sequenceDB kmerDir edgeDir\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <i:kmerDir> <o:edgeDir>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"kmerDir", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::directory },
{"edgeDir", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::directory }}},
{"alignparallel", alignparallel, &par.alignparallel, COMMAND_ALIGNMENT,
"Align candidate edges bucketed by representative key, with many workers",
"# Merge the duplicate copies the k-mer partitions produced, align each pair\n"
"# once, and write the survivors. Workers claim representative-key buckets.\n"
"mmseqs alignparallel sequenceDB edgeDir alnDir --min-seq-id 0.9 -c 0.8\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <i:edgeDir> <o:alnDir>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"edgeDir", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::directory },
{"alnDir", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::directory }}},
{"greedycluster", greedycluster, &par.greedycluster, COMMAND_CLUSTER,
"Greedy clustering over the surviving edges, in one key-ordered sweep",
"# Representatives always have lower keys than their members, so a single\n"
"# left-to-right sweep is the exact greedy. Needs two bits per key, not the\n"
"# eight bytes per sequence stock's fused clustering keeps resident.\n"
"mmseqs greedycluster sequenceDB alnDir clusters.tsv\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <i:alnDir> <o:clusterTsv>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"alnDir", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::directory },
{"clusterTsv", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"mergeclusterparallel", mergeclusterparallel, &par.mergeclusterparallel, COMMAND_CLUSTER,
"Compose the clusterings of successive linclust passes",
"# Folds a later clustering into an earlier one by a key-range join, so no\n"
"# per-sequence list array is needed (stock keeps 24 B/sequence of empty\n"
"# list headers before storing a single member).\n"
"mmseqs mergeclusterparallel sequenceDB pass1.tsv pass2.tsv clusters.tsv\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <i:clusterTsv1> ... <i:clusterTsvN> <o:clusterTsv>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"clusterTsv", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA|DbType::VARIADIC, &DbValidator::flatfile },
{"clusterTsv", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"createrepdb", createrepdb, &par.createrepdb, COMMAND_HIDDEN,
"Build a densely re-keyed representative DB for the next linclust pass",
"# An internal stage of linclustparallel, not a general database builder: the\n"
"# output carries a dense companion index instead of the usual .index/.lookup,\n"
"# so only the parallel linclust stages can read it.\n"
"# Representatives keep length order, so sub-key i is the i-th representative\n"
"# and the copy is sequential. Writes <repDB>.keymap (sub-key -> original key)\n"
"# so the next pass's clustering can be translated back before merging.\n"
"mmseqs createrepdb sequenceDB clusters.tsv repDB\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:sequenceDB> <i:clusterTsv> <o:denseSequenceDB>",
CITATION_MMSEQS2, {{"sequenceDB", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::sequenceDb },
{"clusterTsv", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"sequenceDB", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"translatecluster", translatecluster, &par.translatecluster, COMMAND_CLUSTER,
"Rewrite a clustering from sub-key space into original keys",
"# The pass over a createrepdb sub-database returns sub-keys; this maps them\n"
"# back before merging. Each key column is translated in its own bucketed\n"
"# pass, so the key map is only ever read in contiguous slices.\n"
"mmseqs translatecluster pass2.tsv repDB.keymap pass2_orig.tsv\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:clusterTsv> <i:keyMap> <o:clusterTsv>",
CITATION_MMSEQS2, {{"clusterTsv", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"keyMap", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"clusterTsv", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"translatekeys", translatekeys, &par.translatekeys, COMMAND_CLUSTER,
"Rewrite a clustering from database keys into accessions",
"# The distributed pipeline works in dense keys, so its clusters.tsv holds\n"
"# numbers. This joins it against the database's .lookup by streaming: each\n"
"# key column is translated in its own bucketed pass, and the lookup is read\n"
"# sequentially, so nothing per-key is ever resident.\n"
"mmseqs translatekeys clusters.tsv sequenceDB.lookup clusters_named.tsv\n",
"Björn Buschkämper <bjoern.buschkaemper@gmail.com>",
"<i:clusterTsv> <i:lookupFile> <o:clusterTsv>",
CITATION_MMSEQS2, {{"clusterTsv", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"lookupFile", DbType::ACCESS_MODE_INPUT, DbType::NEED_DATA, &DbValidator::flatfile },
{"clusterTsv", DbType::ACCESS_MODE_OUTPUT, DbType::NEED_DATA, &DbValidator::flatfile }}},
{"kmersearch", kmersearch, &par.kmersearch, COMMAND_PREFILTER,
"Find bottom-m-hashed k-mer matches between target and query DB",
NULL,
Expand Down
25 changes: 0 additions & 25 deletions src/alignment/Align2clust.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -203,31 +203,6 @@ static void pushClusterResult(ClusterResult &&clusterResult) {
}
}

static float parsePrecisionLib(const std::string &scoreFile, double targetSeqid, double targetCov, double targetPrecision) {
std::stringstream in(scoreFile);
std::string line;
int intTargetSeqid = static_cast<int>((targetSeqid + 0.0001) * 100);
int seqIdRest = (intTargetSeqid % 5);
targetSeqid = static_cast<float>(intTargetSeqid - seqIdRest) / 100;
targetCov = static_cast<float>(static_cast<int>((targetCov + 0.0001) * 10)) / 10;

while (std::getline(in, line)) {
std::vector<std::string> values = Util::split(line, " ");
float cov = strtod(values[0].c_str(), NULL);
float seqid = strtod(values[1].c_str(), NULL);
float scorePerCol = strtod(values[2].c_str(), NULL);
float precision = strtod(values[3].c_str(), NULL);
if (MathUtil::AreSame(cov, targetCov) && MathUtil::AreSame(seqid, targetSeqid) && precision >= targetPrecision) {
return scorePerCol;
}
}

Debug(Debug::WARNING) << "Can not find any score per column for coverage "
<< targetCov << " and sequence identity " << targetSeqid
<< ". No hit will be filtered.\n";
return 0;
}

static void writeClustering(DBWriter *dbWriter, const std::pair<DBKeyType, DBKeyType> * results, size_t dbSize) {
std::string resultString;
resultString.reserve(1024 * 1024 * 1024);
Expand Down
29 changes: 28 additions & 1 deletion src/alignment/Matcher.cpp
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
#include <iomanip>
#include <itoa.h>
#include "Matcher.h"
#include <sstream>
#include "MathUtil.h"
#include "Util.h"
#include "Parameters.h"
#include "StripedSmithWaterman.h"
Expand Down Expand Up @@ -411,4 +413,29 @@ std::string getCovSeqidQscPercMinDiagTargetCov() {
reinterpret_cast<const char*>(CovSeqidQscPercMinDiagTargetCov_lib),
CovSeqidQscPercMinDiagTargetCov_lib_len
);
}
}

float parsePrecisionLib(const std::string &scoreFile, double targetSeqid, double targetCov, double targetPrecision) {
std::stringstream in(scoreFile);
std::string line;
int intTargetSeqid = static_cast<int>((targetSeqid + 0.0001) * 100);
int seqIdRest = (intTargetSeqid % 5);
targetSeqid = static_cast<float>(intTargetSeqid - seqIdRest) / 100;
targetCov = static_cast<float>(static_cast<int>((targetCov + 0.0001) * 10)) / 10;

while (std::getline(in, line)) {
std::vector<std::string> values = Util::split(line, " ");
float cov = strtod(values[0].c_str(), NULL);
float seqid = strtod(values[1].c_str(), NULL);
float scorePerCol = strtod(values[2].c_str(), NULL);
float precision = strtod(values[3].c_str(), NULL);
if (MathUtil::AreSame(cov, targetCov) && MathUtil::AreSame(seqid, targetSeqid) && precision >= targetPrecision) {
return scorePerCol;
}
}

Debug(Debug::WARNING) << "Can not find any score per column for coverage "
<< targetCov << " and sequence identity " << targetSeqid
<< ". No hit will be filtered.\n";
return 0;
}
Loading
Loading