-
Notifications
You must be signed in to change notification settings - Fork 3.5k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Wire up DROP CONTINUOUS QUERY #1631
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -3050,6 +3050,17 @@ func (s *Server) CreateContinuousQuery(q *influxql.CreateContinuousQueryStatemen | |
return err | ||
} | ||
|
||
func (s *Server) executeDropContinuousQueryStatement(q *influxql.DropContinuousQueryStatement, user *User) *Result { | ||
return &Result{Err: s.DropContinuousQuery(q)} | ||
} | ||
|
||
// DropContinuousQuery dropsoa continuous query. | ||
func (s *Server) DropContinuousQuery(q *influxql.DropContinuousQueryStatement) error { | ||
c := &dropContinuousQueryCommand{Name: q.Name, Database: q.Database} | ||
_, err := s.broadcast(dropContinuousQueryMessageType, c) | ||
return err | ||
} | ||
|
||
// ContinuousQueries returns a list of all continuous queries. | ||
func (s *Server) ContinuousQueries(database string) []*ContinuousQuery { | ||
s.mu.RLock() | ||
|
@@ -3277,6 +3288,8 @@ func (s *Server) processor(conn MessagingConn, done chan struct{}) { | |
err = s.applySetPrivilege(m) | ||
case createContinuousQueryMessageType: | ||
err = s.applyCreateContinuousQueryCommand(m) | ||
case dropContinuousQueryMessageType: | ||
err = s.applyDropContinuousQueryCommand(m) | ||
case dropSeriesMessageType: | ||
err = s.applyDropSeries(m) | ||
case writeRawSeriesMessageType: | ||
|
@@ -3566,7 +3579,7 @@ func NewContinuousQuery(q string) (*ContinuousQuery, error) { | |
|
||
cq, ok := stmt.(*influxql.CreateContinuousQueryStatement) | ||
if !ok { | ||
return nil, errors.New("query isn't a continuous query") | ||
return nil, errors.New("query isn't a valie continuous query") | ||
} | ||
|
||
cquery := &ContinuousQuery{ | ||
|
@@ -3636,6 +3649,44 @@ func (s *Server) applyCreateContinuousQueryCommand(m *messaging.Message) error { | |
return nil | ||
} | ||
|
||
// applyDropContinuousQueryCommand removes the continuous query from the database object and saves it to the metastore | ||
func (s *Server) applyDropContinuousQueryCommand(m *messaging.Message) error { | ||
var c dropContinuousQueryCommand | ||
|
||
mustUnmarshalJSON(m.Data, &c) | ||
|
||
// retrieve the database and ensure that it exists | ||
db := s.databases[c.Database] | ||
if db == nil { | ||
return ErrDatabaseNotFound | ||
} | ||
|
||
// loop through continuous queries and find the match | ||
cqIndex := -1 | ||
for n, continuousQuery := range db.continuousQueries { | ||
if continuousQuery.cq.Name == c.Name { | ||
cqIndex = n | ||
break | ||
} | ||
} | ||
|
||
if cqIndex == -1 { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. So actually, this could happen, even if everything is OK. If two "drop CQ" commands for the same CQ hit the cluster at different nodes at the same time, one will win, but the other may still have gotten through before the first deletion finished. So if this happens, simply return nil here. In other words, make this operation idempotent. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I guess this could happen in lots of other places too, actually, so maybe it doesn't matter. Worse thing that happens is that the user gets "does not exist" back, which is what he wanted anyway. So doesn't matter. |
||
return ErrContinuousQueryNotFound | ||
} | ||
|
||
// delete the relevant continuous query | ||
copy(db.continuousQueries[cqIndex:], db.continuousQueries[cqIndex+1:]) | ||
db.continuousQueries[len(db.continuousQueries)-1] = nil | ||
db.continuousQueries = db.continuousQueries[:len(db.continuousQueries)-1] | ||
|
||
// persist to metastore | ||
s.meta.mustUpdate(m.Index, func(tx *metatx) error { | ||
return tx.saveDatabase(db) | ||
}) | ||
|
||
return nil | ||
} | ||
|
||
// RunContinuousQueries will run any continuous queries that are due to run and write the | ||
// results back into the database | ||
func (s *Server) RunContinuousQueries() error { | ||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
You should check here if the database and CQ actually exist.