mirror of
https://github.com/metabarcoding/obitools4.git
synced 2026-08-24 13:51:18 +00:00
feat: add JSON input support for biological sequences
Introduce a streaming JSON parser that decodes biological sequences using `goccy/go-json` with configurable batching to minimize memory overhead. Extend the CLI, file suffix filters, and MIME type detection to automatically recognize and route JSON inputs. Refactor header parsing into a centralized switch-case handler for improved maintainability.
This commit is contained in:
@@ -199,8 +199,111 @@ func _parse_json_array_interface(str []byte) ([]interface{}, error) {
|
||||
return values, nil
|
||||
}
|
||||
|
||||
func _parse_json_header_(header string, sequence *obiseq.BioSequence) string {
|
||||
// _parse_json_annotation_field parses a single key/value pair coming from a
|
||||
// JSON object (either a FASTA/FASTQ inline JSON header, or the "annotations"
|
||||
// field of a JSON sequence record) and applies it to the sequence, special
|
||||
// casing the well-known OBITools attributes (id, definition, count, taxid,
|
||||
// obiclean_*, merged_*).
|
||||
func _parse_json_annotation_field(key []byte, value []byte, dataType jsonparser.ValueType, sequence *obiseq.BioSequence) error {
|
||||
annotations := sequence.Annotations()
|
||||
var err error
|
||||
|
||||
skey := obiutils.UnsafeString(key)
|
||||
|
||||
switch {
|
||||
case skey == "id":
|
||||
sequence.SetId(string(value))
|
||||
case skey == "definition":
|
||||
sequence.SetDefinition(string(value))
|
||||
|
||||
case skey == "count":
|
||||
if dataType != jsonparser.Number {
|
||||
log.Fatalf("%s: Count attribut must be numeric: %s", sequence.Id(), string(value))
|
||||
}
|
||||
count, err := jsonparser.ParseInt(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse count %s", sequence.Id(), string(value))
|
||||
}
|
||||
sequence.SetCount(int(count))
|
||||
|
||||
case skey == "obiclean_weight":
|
||||
weight, err := _parse_json_map_int(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse obiclean weight %s", sequence.Id(), string(value))
|
||||
}
|
||||
annotations[skey] = weight
|
||||
|
||||
case skey == "obiclean_status":
|
||||
status, err := _parse_json_map_string(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse obiclean status %s", sequence.Id(), string(value))
|
||||
}
|
||||
annotations[skey] = status
|
||||
|
||||
case strings.HasPrefix(skey, "merged_"):
|
||||
if dataType == jsonparser.Object {
|
||||
data, err := _parse_json_map_int(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse merged slot %s: %v", sequence.Id(), skey, err)
|
||||
} else {
|
||||
annotations[skey] = obiseq.MapAsStatsOnValues(data)
|
||||
}
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse merged slot %s", sequence.Id(), skey)
|
||||
}
|
||||
|
||||
case skey == "taxid":
|
||||
if dataType == jsonparser.Number || dataType == jsonparser.String {
|
||||
taxid := string(value)
|
||||
sequence.SetTaxid(taxid)
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse taxid %s", sequence.Id(), string(value))
|
||||
}
|
||||
|
||||
case strings.HasSuffix(skey, "_taxid"):
|
||||
if dataType == jsonparser.Number || dataType == jsonparser.String {
|
||||
rank := skey[:len(skey)-len("_taxid")]
|
||||
|
||||
taxid := string(value)
|
||||
sequence.SetTaxid(taxid, rank)
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse taxid %s", sequence.Id(), string(value))
|
||||
}
|
||||
|
||||
default:
|
||||
skey = strings.Clone(skey)
|
||||
switch dataType {
|
||||
case jsonparser.String:
|
||||
annotations[skey] = string(value)
|
||||
case jsonparser.Number:
|
||||
// Try to parse the number as an int at first then as float if that fails.
|
||||
annotations[skey], err = jsonparser.ParseInt(value)
|
||||
if err != nil {
|
||||
annotations[skey], err = strconv.ParseFloat(obiutils.UnsafeString(value), 64)
|
||||
}
|
||||
case jsonparser.Array:
|
||||
annotations[skey], err = _parse_json_array_interface(value)
|
||||
case jsonparser.Object:
|
||||
annotations[skey], err = _parse_json_map_interface(value)
|
||||
case jsonparser.Boolean:
|
||||
annotations[skey], err = jsonparser.ParseBoolean(value)
|
||||
case jsonparser.Null:
|
||||
annotations[skey] = nil
|
||||
default:
|
||||
log.Fatalf("Unknown data type %v", dataType)
|
||||
}
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
annotations[skey] = "NaN"
|
||||
log.Fatalf("%s: Cannot parse value %s assicated to key %s into a %s value",
|
||||
sequence.Id(), string(value), skey, dataType.String())
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
func _parse_json_header_(header string, sequence *obiseq.BioSequence) string {
|
||||
start := -1
|
||||
stop := -1
|
||||
level := 0
|
||||
@@ -240,101 +343,7 @@ func _parse_json_header_(header string, sequence *obiseq.BioSequence) string {
|
||||
|
||||
jsonparser.ObjectEach(obiutils.UnsafeBytes(header[start:stop]),
|
||||
func(key []byte, value []byte, dataType jsonparser.ValueType, offset int) error {
|
||||
var err error
|
||||
|
||||
skey := obiutils.UnsafeString(key)
|
||||
|
||||
switch {
|
||||
case skey == "id":
|
||||
sequence.SetId(string(value))
|
||||
case skey == "definition":
|
||||
sequence.SetDefinition(string(value))
|
||||
|
||||
case skey == "count":
|
||||
if dataType != jsonparser.Number {
|
||||
log.Fatalf("%s: Count attribut must be numeric: %s", sequence.Id(), string(value))
|
||||
}
|
||||
count, err := jsonparser.ParseInt(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse count %s", sequence.Id(), string(value))
|
||||
}
|
||||
sequence.SetCount(int(count))
|
||||
|
||||
case skey == "obiclean_weight":
|
||||
weight, err := _parse_json_map_int(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse obiclean weight %s", sequence.Id(), string(value))
|
||||
}
|
||||
annotations[skey] = weight
|
||||
|
||||
case skey == "obiclean_status":
|
||||
status, err := _parse_json_map_string(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse obiclean status %s", sequence.Id(), string(value))
|
||||
}
|
||||
annotations[skey] = status
|
||||
|
||||
case strings.HasPrefix(skey, "merged_"):
|
||||
if dataType == jsonparser.Object {
|
||||
data, err := _parse_json_map_int(value)
|
||||
if err != nil {
|
||||
log.Fatalf("%s: Cannot parse merged slot %s: %v", sequence.Id(), skey, err)
|
||||
} else {
|
||||
annotations[skey] = obiseq.MapAsStatsOnValues(data)
|
||||
}
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse merged slot %s", sequence.Id(), skey)
|
||||
}
|
||||
|
||||
case skey == "taxid":
|
||||
if dataType == jsonparser.Number || dataType == jsonparser.String {
|
||||
taxid := string(value)
|
||||
sequence.SetTaxid(taxid)
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse taxid %s", sequence.Id(), string(value))
|
||||
}
|
||||
|
||||
case strings.HasSuffix(skey, "_taxid"):
|
||||
if dataType == jsonparser.Number || dataType == jsonparser.String {
|
||||
rank := skey[:len(skey)-len("_taxid")]
|
||||
|
||||
taxid := string(value)
|
||||
sequence.SetTaxid(taxid, rank)
|
||||
} else {
|
||||
log.Fatalf("%s: Cannot parse taxid %s", sequence.Id(), string(value))
|
||||
}
|
||||
|
||||
default:
|
||||
skey = strings.Clone(skey)
|
||||
switch dataType {
|
||||
case jsonparser.String:
|
||||
annotations[skey] = string(value)
|
||||
case jsonparser.Number:
|
||||
// Try to parse the number as an int at first then as float if that fails.
|
||||
annotations[skey], err = jsonparser.ParseInt(value)
|
||||
if err != nil {
|
||||
annotations[skey], err = strconv.ParseFloat(obiutils.UnsafeString(value), 64)
|
||||
}
|
||||
case jsonparser.Array:
|
||||
annotations[skey], err = _parse_json_array_interface(value)
|
||||
case jsonparser.Object:
|
||||
annotations[skey], err = _parse_json_map_interface(value)
|
||||
case jsonparser.Boolean:
|
||||
annotations[skey], err = jsonparser.ParseBoolean(value)
|
||||
case jsonparser.Null:
|
||||
annotations[skey] = nil
|
||||
default:
|
||||
log.Fatalf("Unknown data type %v", dataType)
|
||||
}
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
annotations[skey] = "NaN"
|
||||
log.Fatalf("%s: Cannot parse value %s assicated to key %s into a %s value",
|
||||
sequence.Id(), string(value), skey, dataType.String())
|
||||
}
|
||||
|
||||
return err
|
||||
return _parse_json_annotation_field(key, value, dataType, sequence)
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
package obiformats
|
||||
|
||||
import (
|
||||
"io"
|
||||
"os"
|
||||
"path"
|
||||
|
||||
"git.metabarcoding.org/obitools/obitools4/obitools4/pkg/obidefault"
|
||||
"git.metabarcoding.org/obitools/obitools4/obitools4/pkg/obiiter"
|
||||
"git.metabarcoding.org/obitools/obitools4/obitools4/pkg/obiseq"
|
||||
"git.metabarcoding.org/obitools/obitools4/obitools4/pkg/obiutils"
|
||||
"github.com/buger/jsonparser"
|
||||
"github.com/goccy/go-json"
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
// _parse_json_record parses a single JSON object describing a sequence
|
||||
// (as produced by JSONRecord in json_writer.go) into a *obiseq.BioSequence.
|
||||
func _parse_json_record(raw []byte, shift byte) *obiseq.BioSequence {
|
||||
sequence := obiseq.NewEmptyBioSequence(0)
|
||||
|
||||
if id, err := jsonparser.GetString(raw, "id"); err == nil {
|
||||
sequence.SetId(id)
|
||||
}
|
||||
|
||||
if seq, err := jsonparser.GetString(raw, "sequence"); err == nil {
|
||||
sequence.SetSequence([]byte(seq))
|
||||
}
|
||||
|
||||
if qual, err := jsonparser.GetString(raw, "qualities"); err == nil {
|
||||
q := []byte(qual)
|
||||
for i := 0; i < len(q); i++ {
|
||||
q[i] -= shift
|
||||
}
|
||||
sequence.SetQualities(q)
|
||||
}
|
||||
|
||||
if annot, dataType, _, err := jsonparser.Get(raw, "annotations"); err == nil && dataType == jsonparser.Object {
|
||||
jsonparser.ObjectEach(annot,
|
||||
func(key []byte, value []byte, valType jsonparser.ValueType, offset int) error {
|
||||
return _parse_json_annotation_field(key, value, valType, sequence)
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
return sequence
|
||||
}
|
||||
|
||||
// _ParseJsonFile streams the top-level JSON array, decoding and pushing one
|
||||
// batch of sequences at a time, without ever loading the whole document in
|
||||
// memory. Only one raw record at a time is buffered by the decoder.
|
||||
func _ParseJsonFile(source string,
|
||||
reader io.Reader,
|
||||
out obiiter.IBioSequence,
|
||||
shift byte,
|
||||
batchSize int) {
|
||||
|
||||
dec := json.NewDecoder(reader)
|
||||
|
||||
if _, err := dec.Token(); err != nil {
|
||||
if err == io.EOF {
|
||||
out.Done()
|
||||
return
|
||||
}
|
||||
log.Fatalf("cannot parse JSON data: %v", err)
|
||||
}
|
||||
|
||||
slice := obiseq.MakeBioSequenceSlice()
|
||||
o := 0
|
||||
|
||||
for dec.More() {
|
||||
var raw json.RawMessage
|
||||
|
||||
if err := dec.Decode(&raw); err != nil {
|
||||
log.Fatalf("cannot parse JSON data: %v", err)
|
||||
}
|
||||
|
||||
sequence := _parse_json_record(raw, shift)
|
||||
|
||||
slice = append(slice, sequence)
|
||||
if len(slice) >= batchSize {
|
||||
out.Push(obiiter.MakeBioSequenceBatch(source, o, slice))
|
||||
o++
|
||||
slice = obiseq.MakeBioSequenceSlice()
|
||||
}
|
||||
}
|
||||
|
||||
if len(slice) > 0 {
|
||||
out.Push(obiiter.MakeBioSequenceBatch(source, o, slice))
|
||||
}
|
||||
|
||||
out.Done()
|
||||
}
|
||||
|
||||
func ReadJSON(reader io.Reader, options ...WithOption) (obiiter.IBioSequence, error) {
|
||||
|
||||
opt := MakeOptions(options)
|
||||
out := obiiter.MakeIBioSequence()
|
||||
|
||||
out.Add(1)
|
||||
go _ParseJsonFile(opt.Source(),
|
||||
reader,
|
||||
out,
|
||||
obidefault.ReadQualitiesShift(),
|
||||
opt.BatchSize())
|
||||
|
||||
go func() {
|
||||
out.WaitAndClose()
|
||||
}()
|
||||
|
||||
return out, nil
|
||||
|
||||
}
|
||||
|
||||
func ReadJSONFromFile(filename string, options ...WithOption) (obiiter.IBioSequence, error) {
|
||||
|
||||
options = append(options, OptionsSource(obiutils.RemoveAllExt((path.Base(filename)))))
|
||||
file, err := obiutils.Ropen(filename)
|
||||
|
||||
if err == obiutils.ErrNoContent {
|
||||
log.Infof("file %s is empty", filename)
|
||||
return ReadEmptyFile(options...)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return obiiter.NilIBioSequence, err
|
||||
}
|
||||
|
||||
return ReadJSON(file, options...)
|
||||
}
|
||||
|
||||
func ReadJSONFromStdin(reader io.Reader, options ...WithOption) (obiiter.IBioSequence, error) {
|
||||
options = append(options, OptionsSource(obiutils.RemoveAllExt("stdin")))
|
||||
input, err := obiutils.Buf(os.Stdin)
|
||||
|
||||
if err == obiutils.ErrNoContent {
|
||||
log.Infof("stdin is empty")
|
||||
return ReadEmptyFile(options...)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
log.Fatalf("open file error: %v", err)
|
||||
return obiiter.NilIBioSequence, err
|
||||
}
|
||||
|
||||
return ReadJSON(input, options...)
|
||||
}
|
||||
@@ -145,6 +145,8 @@ func ReadSequencesFromFile(filename string,
|
||||
return ReadGenbank(reader, options...)
|
||||
case "text/csv":
|
||||
return ReadCSV(reader, options...)
|
||||
case "application/json":
|
||||
return ReadJSON(reader, options...)
|
||||
default:
|
||||
log.Fatalf("File %s has guessed format %s which is not yet implemented",
|
||||
filename, mime.String())
|
||||
|
||||
Reference in New Issue
Block a user