2019-10-03 03:06:03 +08:00
|
|
|
package weed_server
|
|
|
|
|
|
|
|
import (
|
2022-07-29 15:17:28 +08:00
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/operation"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/query/json"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
2019-10-03 03:06:03 +08:00
|
|
|
"github.com/tidwall/gjson"
|
|
|
|
)
|
|
|
|
|
|
|
|
func (vs *VolumeServer) Query(req *volume_server_pb.QueryRequest, stream volume_server_pb.VolumeServer_QueryServer) error {
|
|
|
|
|
|
|
|
for _, fid := range req.FromFileIds {
|
|
|
|
|
|
|
|
vid, id_cookie, err := operation.ParseFileId(fid)
|
|
|
|
if err != nil {
|
|
|
|
glog.V(0).Infof("volume query failed to parse fid %s: %v", fid, err)
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
n := new(needle.Needle)
|
|
|
|
volumeId, _ := needle.NewVolumeId(vid)
|
|
|
|
n.ParsePath(id_cookie)
|
|
|
|
|
|
|
|
cookie := n.Cookie
|
2021-08-09 14:25:16 +08:00
|
|
|
if _, err := vs.store.ReadVolumeNeedle(volumeId, n, nil, nil); err != nil {
|
2019-10-03 03:06:03 +08:00
|
|
|
glog.V(0).Infof("volume query failed to read fid %s: %v", fid, err)
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if n.Cookie != cookie {
|
|
|
|
glog.V(0).Infof("volume query failed to read fid cookie %s: %v", fid, err)
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2019-10-09 15:03:18 +08:00
|
|
|
if req.InputSerialization.CsvInput != nil {
|
2019-10-03 03:06:03 +08:00
|
|
|
|
|
|
|
}
|
|
|
|
|
2019-10-09 15:03:18 +08:00
|
|
|
if req.InputSerialization.JsonInput != nil {
|
2019-10-03 03:06:03 +08:00
|
|
|
|
2019-10-07 13:35:05 +08:00
|
|
|
stripe := &volume_server_pb.QueriedStripe{
|
2019-10-09 15:03:18 +08:00
|
|
|
Records: nil,
|
2019-10-07 13:35:05 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
filter := json.Query{
|
|
|
|
Field: req.Filter.Field,
|
2019-10-09 15:03:18 +08:00
|
|
|
Op: req.Filter.Operand,
|
|
|
|
Value: req.Filter.Value,
|
2019-10-03 03:06:03 +08:00
|
|
|
}
|
2019-10-09 15:03:18 +08:00
|
|
|
gjson.ForEachLine(string(n.Data), func(line gjson.Result) bool {
|
2019-10-07 13:35:05 +08:00
|
|
|
passedFilter, values := json.QueryJson(line.Raw, req.Selections, filter)
|
|
|
|
if !passedFilter {
|
|
|
|
return true
|
|
|
|
}
|
|
|
|
stripe.Records = json.ToJson(stripe.Records, req.Selections, values)
|
2019-10-03 03:06:03 +08:00
|
|
|
return true
|
|
|
|
})
|
2019-10-07 13:35:05 +08:00
|
|
|
err = stream.Send(stripe)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2019-10-03 03:06:03 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
2019-10-09 15:03:18 +08:00
|
|
|
}
|