This is an automated email from the ASF dual-hosted git repository.

jackietien pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iotdb-client-go.git


The following commit(s) were added to refs/heads/main by this push:
     new bbff3ce  Add FastLastDataQueryForOnePrefixPath interface
bbff3ce is described below

commit bbff3ce5d77194f1a1f61910bfed922d603ed2a6
Author: shuwenwei <[email protected]>
AuthorDate: Sun Jun 22 16:09:04 2025 +0800

    Add FastLastDataQueryForOnePrefixPath interface
---
 client/session.go | 54 ++++++++++++++++++++++++++++++++++++++++++++++++++++++
 1 file changed, 54 insertions(+)

diff --git a/client/session.go b/client/session.go
index 488d04c..230799f 100644
--- a/client/session.go
+++ b/client/session.go
@@ -601,6 +601,60 @@ func (s *Session) 
ExecuteAggregationQueryWithLegalNodes(paths []string, aggregat
        }
 }
 
+func (s *Session) ExecuteGroupByQueryIntervalQuery(database *string, device, 
measurement string, aggregationType common.TAggregationType,
+       dataType int32, startTime *int64, endTime *int64, interval *int64, 
timeoutMs *int64, isAligned *bool) (*SessionDataSet, error) {
+
+       request := rpc.TSGroupByQueryIntervalReq{SessionId: s.sessionId, 
StatementId: s.requestStatementId,
+               Database: database, Device: device, Measurement: measurement, 
AggregationType: aggregationType, DataType: dataType,
+               StartTime: startTime, EndTime: endTime, Interval: interval, 
FetchSize: &s.config.FetchSize,
+               Timeout: timeoutMs, IsAligned: isAligned}
+       if resp, err := 
s.client.ExecuteGroupByQueryIntervalQuery(context.Background(), &request); err 
== nil {
+               if statusErr := VerifySuccess(resp.Status); statusErr == nil {
+                       return NewSessionDataSet("", resp.Columns, 
resp.DataTypeList, resp.ColumnNameIndexMap, *resp.QueryId, 
s.requestStatementId, s.client, s.sessionId, resp.QueryResult_, 
resp.IgnoreTimeStamp != nil && *resp.IgnoreTimeStamp, timeoutMs, 
*resp.MoreData, s.config.FetchSize, s.config.TimeZone, s.timeFactor, 
resp.GetColumnIndex2TsBlockColumnIndexList())
+               } else {
+                       return nil, statusErr
+               }
+       } else {
+               if s.reconnect() {
+                       request.SessionId = s.sessionId
+                       resp, err = 
s.client.ExecuteGroupByQueryIntervalQuery(context.Background(), &request)
+                       if statusErr := VerifySuccess(resp.Status); statusErr 
== nil {
+                               return NewSessionDataSet("", resp.Columns, 
resp.DataTypeList, resp.ColumnNameIndexMap, *resp.QueryId, 
s.requestStatementId, s.client, s.sessionId, resp.QueryResult_, 
resp.IgnoreTimeStamp != nil && *resp.IgnoreTimeStamp, timeoutMs, 
*resp.MoreData, s.config.FetchSize, s.config.TimeZone, s.timeFactor, 
resp.GetColumnIndex2TsBlockColumnIndexList())
+                       } else {
+                               return nil, statusErr
+                       }
+               }
+               return nil, err
+       }
+}
+
+func (s *Session) ExecuteFastLastDataQueryForOnePrefixPath(prefixes []string, 
timeoutMs *int64) (*SessionDataSet, error) {
+       request := rpc.TSFastLastDataQueryForOnePrefixPathReq{
+               SessionId:   s.sessionId,
+               StatementId: s.requestStatementId,
+               Prefixes:    prefixes,
+               Timeout:     timeoutMs,
+       }
+       if resp, err := 
s.client.ExecuteFastLastDataQueryForOnePrefixPath(context.Background(), 
&request); err == nil {
+               if statusErr := VerifySuccess(resp.Status); statusErr == nil {
+                       return NewSessionDataSet("", resp.Columns, 
resp.DataTypeList, resp.ColumnNameIndexMap, *resp.QueryId, 
s.requestStatementId, s.client, s.sessionId, resp.QueryResult_, 
resp.IgnoreTimeStamp != nil && *resp.IgnoreTimeStamp, timeoutMs, 
*resp.MoreData, s.config.FetchSize, s.config.TimeZone, s.timeFactor, 
resp.GetColumnIndex2TsBlockColumnIndexList())
+               } else {
+                       return nil, statusErr
+               }
+       } else {
+               if s.reconnect() {
+                       request.SessionId = s.sessionId
+                       resp, err = 
s.client.ExecuteFastLastDataQueryForOnePrefixPath(context.Background(), 
&request)
+                       if statusErr := VerifySuccess(resp.Status); statusErr 
== nil {
+                               return NewSessionDataSet("", resp.Columns, 
resp.DataTypeList, resp.ColumnNameIndexMap, *resp.QueryId, 
s.requestStatementId, s.client, s.sessionId, resp.QueryResult_, 
resp.IgnoreTimeStamp != nil && *resp.IgnoreTimeStamp, timeoutMs, 
*resp.MoreData, s.config.FetchSize, s.config.TimeZone, s.timeFactor, 
resp.GetColumnIndex2TsBlockColumnIndexList())
+                       } else {
+                               return nil, statusErr
+                       }
+               }
+               return nil, err
+       }
+}
+
 func (s *Session) genTSInsertRecordReq(deviceId string, time int64,
        measurements []string,
        types []TSDataType,

Reply via email to