Skip to content
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

Calculate topology based on registered data nodes #324

Closed
wants to merge 9 commits into from
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 6 additions & 1 deletion banyand/liaison/grpc/discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,17 +41,19 @@ var errNotExist = errors.New("the object doesn't exist")

type discoveryService struct {
pipeline queue.Client
nodeRegistry NodeRegistry
metadataRepo metadata.Repo
shardRepo *shardRepo
entityRepo *entityRepo
log *logger.Logger
kind schema.Kind
}

func newDiscoveryService(pipeline queue.Client, kind schema.Kind, metadataRepo metadata.Repo) *discoveryService {
func newDiscoveryService(pipeline queue.Client, kind schema.Kind, metadataRepo metadata.Repo, nodeRegistry NodeRegistry) *discoveryService {
sr := &shardRepo{shardEventsMap: make(map[identity]uint32)}
er := &entityRepo{entitiesMap: make(map[identity]partition.EntityLocator)}
return &discoveryService{
nodeRegistry: nodeRegistry,
shardRepo: sr,
entityRepo: er,
pipeline: pipeline,
Expand All @@ -61,6 +63,9 @@ func newDiscoveryService(pipeline queue.Client, kind schema.Kind, metadataRepo m
}

func (ds *discoveryService) initialize(ctx context.Context) error {
if err := ds.nodeRegistry.Initialize(); err != nil {
return err
}
ctxLocal, cancel := context.WithTimeout(ctx, 5*time.Second)
groups, err := ds.metadataRepo.GroupRegistry().ListGroup(ctxLocal)
cancel()
Expand Down
13 changes: 10 additions & 3 deletions banyand/liaison/grpc/measure.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,11 +113,18 @@ func (ms *measureService) Write(measure measurev1.MeasureService_WriteServer) er
SeriesHash: tsdb.HashEntity(entity),
EntityValues: tagValues.Encode(),
}
// TODO: set node id
message := bus.NewMessageWithNode(bus.MessageID(time.Now().UnixNano()), "todo", iwr)
nodeID, errPickNode := ms.nodeRegistry.Locate(writeRequest.GetMetadata().GetGroup(), writeRequest.GetMetadata().GetName(), uint32(shardID))
if errPickNode != nil {
ms.sampled.Error().Err(errPickNode).RawJSON("written", logger.Proto(writeRequest)).Msg("failed to pick an available node")
reply(writeRequest.GetMetadata(), modelv1.Status_STATUS_INTERNAL_ERROR, writeRequest.GetMessageId(), measure, ms.sampled)
continue
}
message := bus.NewMessageWithNode(bus.MessageID(time.Now().UnixNano()), nodeID, iwr)
_, errWritePub := publisher.Publish(data.TopicMeasureWrite, message)
if errWritePub != nil {
ms.sampled.Error().Err(errWritePub).RawJSON("written", logger.Proto(writeRequest)).Msg("failed to send a message")
ms.sampled.Error().Err(errWritePub).RawJSON("written", logger.Proto(writeRequest)).Str("nodeID", nodeID).Msg("failed to send a message")
reply(writeRequest.GetMetadata(), modelv1.Status_STATUS_INTERNAL_ERROR, writeRequest.GetMessageId(), measure, ms.sampled)
continue
}
reply(nil, modelv1.Status_STATUS_SUCCEED, writeRequest.GetMessageId(), measure, ms.sampled)
}
Expand Down
110 changes: 110 additions & 0 deletions banyand/liaison/grpc/node.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
// Licensed to Apache Software Foundation (ASF) under one or more contributor
// license agreements. See the NOTICE file distributed with
// this work for additional information regarding copyright
// ownership. Apache Software Foundation (ASF) licenses this file to you under
// the Apache License, Version 2.0 (the "License"); you may
// not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

package grpc

import (
"sync"

"github.com/pkg/errors"

databasev1 "github.com/apache/skywalking-banyandb/api/proto/banyandb/database/v1"
"github.com/apache/skywalking-banyandb/banyand/metadata"
"github.com/apache/skywalking-banyandb/banyand/metadata/schema"
"github.com/apache/skywalking-banyandb/pkg/node"
)

var (
_ schema.EventHandler = (*clusterNodeService)(nil)
_ NodeRegistry = (*clusterNodeService)(nil)
)

// NodeRegistry is for locating data node with group/name of the metadata
// together with the shardID calculated from the incoming data.
type NodeRegistry interface {
Initialize() error
Locate(group, name string, shardID uint32) (string, error)
}

type clusterNodeService struct {
metaRepo metadata.Repo
sel node.Selector
sync.Once
}

// NewClusterNodeRegistry creates a cluster-aware node registry.
func NewClusterNodeRegistry(metaRepo metadata.Repo, selector node.Selector) NodeRegistry {
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

#333 will use queue.Client to listen to node events, reducing the need for the liaison node to set up two etcd watchers. Please merge this once PR is merged.

return &clusterNodeService{
metaRepo: metaRepo,
sel: selector,
}
}

func (n *clusterNodeService) Initialize() error {
n.Do(func() {
n.metaRepo.RegisterHandler("cluster-node-service", schema.KindNode, n)
})
return nil
}

func (n *clusterNodeService) Locate(group, name string, shardID uint32) (string, error) {
nodeID, err := n.sel.Pick(group, name, shardID)
if err != nil {
return "", errors.Wrapf(err, "fail to locate %s/%s(%d)", group, name, shardID)
}
return nodeID, nil
}

func (n *clusterNodeService) OnAddOrUpdate(metadata schema.Metadata) {
switch metadata.Kind {
case schema.KindNode:
inputNode := metadata.Spec.(*databasev1.Node)
if inputNode.Metadata.GetName() == "" {
return
}
n.sel.AddNode(inputNode)
default:
}
}

func (n *clusterNodeService) OnDelete(metadata schema.Metadata) {
switch metadata.Kind {
case schema.KindNode:
dNode := metadata.Spec.(*databasev1.Node)
if dNode.Metadata.GetName() == "" {
return
}
n.sel.RemoveNode(dNode)
default:
}
}

type localNodeService struct{}

// NewLocalNodeRegistry creates a local(fake) node registry.
func NewLocalNodeRegistry() NodeRegistry {
return localNodeService{}
}

func (localNodeService) Initialize() error {
return nil
}

// Locate of localNodeService always returns local.
func (localNodeService) Locate(_, _ string, _ uint32) (string, error) {
return "local", nil
}
2 changes: 1 addition & 1 deletion banyand/liaison/grpc/registry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ func setupForRegistry() func() {
metaSvc, err := metadata.NewService(context.TODO())
Expect(err).NotTo(HaveOccurred())

tcp := grpc.NewServer(context.TODO(), pipeline, metaSvc)
tcp := grpc.NewServer(context.TODO(), pipeline, metaSvc, grpc.NewLocalNodeRegistry())
preloadStreamSvc := &preloadStreamService{metaSvc: metaSvc}
var flags []string
metaPath, metaDeferFunc, err := test.NewSpace()
Expand Down
6 changes: 3 additions & 3 deletions banyand/liaison/grpc/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,12 +91,12 @@ type server struct {
}

// NewServer returns a new gRPC server.
func NewServer(_ context.Context, pipeline queue.Client, schemaRegistry metadata.Repo) Server {
func NewServer(_ context.Context, pipeline queue.Client, schemaRegistry metadata.Repo, nodeRegistry NodeRegistry) Server {
streamSVC := &streamService{
discoveryService: newDiscoveryService(pipeline, schema.KindStream, schemaRegistry),
discoveryService: newDiscoveryService(pipeline, schema.KindStream, schemaRegistry, nodeRegistry),
}
measureSVC := &measureService{
discoveryService: newDiscoveryService(pipeline, schema.KindMeasure, schemaRegistry),
discoveryService: newDiscoveryService(pipeline, schema.KindMeasure, schemaRegistry, nodeRegistry),
}
s := &server{
pipeline: pipeline,
Expand Down
12 changes: 10 additions & 2 deletions banyand/liaison/grpc/stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,10 +115,18 @@ func (s *streamService) Write(stream streamv1.StreamService_WriteServer) error {
if s.log.Debug().Enabled() {
iwr.EntityValues = tagValues.Encode()
}
message := bus.NewMessageWithNode(bus.MessageID(time.Now().UnixNano()), "TODO", iwr)
nodeID, errPickNode := s.nodeRegistry.Locate(writeEntity.GetMetadata().GetGroup(), writeEntity.GetMetadata().GetName(), uint32(shardID))
if errPickNode != nil {
s.sampled.Error().Err(errPickNode).RawJSON("written", logger.Proto(writeEntity)).Msg("failed to pick an available node")
reply(writeEntity.GetMetadata(), modelv1.Status_STATUS_INTERNAL_ERROR, writeEntity.GetMessageId(), stream, s.sampled)
continue
}
message := bus.NewMessageWithNode(bus.MessageID(time.Now().UnixNano()), nodeID, iwr)
_, errWritePub := publisher.Publish(data.TopicStreamWrite, message)
if errWritePub != nil {
s.sampled.Error().Err(errWritePub).RawJSON("written", logger.Proto(writeEntity)).Msg("failed to send a message")
s.sampled.Error().Err(errWritePub).RawJSON("written", logger.Proto(writeEntity)).Str("nodeID", nodeID).Msg("failed to send a message")
reply(writeEntity.GetMetadata(), modelv1.Status_STATUS_INTERNAL_ERROR, writeEntity.GetMessageId(), stream, s.sampled)
continue
}
reply(nil, modelv1.Status_STATUS_SUCCEED, writeEntity.GetMessageId(), stream, s.sampled)
}
Expand Down
7 changes: 7 additions & 0 deletions dist/LICENSE
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,7 @@ Apache-2.0 licenses
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0 Apache-2.0
github.com/inconshreveable/mousetrap v1.1.0 Apache-2.0
github.com/jonboulle/clockwork v0.4.0 Apache-2.0
github.com/kkdai/maglev v0.2.0 Apache-2.0
github.com/matttproud/golang_protobuf_extensions v1.0.4 Apache-2.0
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd Apache-2.0
github.com/modern-go/reflect2 v1.0.2 Apache-2.0
Expand Down Expand Up @@ -296,6 +297,12 @@ BSD-3-Clause and Apache-2.0 and MIT licenses

github.com/klauspost/compress v1.15.6 BSD-3-Clause and Apache-2.0 and MIT

========================================================================
CC0-1.0 licenses
========================================================================

github.com/dchest/siphash v1.2.3 CC0-1.0

========================================================================
ISC licenses
========================================================================
Expand Down
121 changes: 121 additions & 0 deletions dist/licenses/license-github.com-dchest-siphash.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
Creative Commons Legal Code

CC0 1.0 Universal

CREATIVE COMMONS CORPORATION IS NOT A LAW FIRM AND DOES NOT PROVIDE
LEGAL SERVICES. DISTRIBUTION OF THIS DOCUMENT DOES NOT CREATE AN
ATTORNEY-CLIENT RELATIONSHIP. CREATIVE COMMONS PROVIDES THIS
INFORMATION ON AN "AS-IS" BASIS. CREATIVE COMMONS MAKES NO WARRANTIES
REGARDING THE USE OF THIS DOCUMENT OR THE INFORMATION OR WORKS
PROVIDED HEREUNDER, AND DISCLAIMS LIABILITY FOR DAMAGES RESULTING FROM
THE USE OF THIS DOCUMENT OR THE INFORMATION OR WORKS PROVIDED
HEREUNDER.

Statement of Purpose

The laws of most jurisdictions throughout the world automatically confer
exclusive Copyright and Related Rights (defined below) upon the creator
and subsequent owner(s) (each and all, an "owner") of an original work of
authorship and/or a database (each, a "Work").

Certain owners wish to permanently relinquish those rights to a Work for
the purpose of contributing to a commons of creative, cultural and
scientific works ("Commons") that the public can reliably and without fear
of later claims of infringement build upon, modify, incorporate in other
works, reuse and redistribute as freely as possible in any form whatsoever
and for any purposes, including without limitation commercial purposes.
These owners may contribute to the Commons to promote the ideal of a free
culture and the further production of creative, cultural and scientific
works, or to gain reputation or greater distribution for their Work in
part through the use and efforts of others.

For these and/or other purposes and motivations, and without any
expectation of additional consideration or compensation, the person
associating CC0 with a Work (the "Affirmer"), to the extent that he or she
is an owner of Copyright and Related Rights in the Work, voluntarily
elects to apply CC0 to the Work and publicly distribute the Work under its
terms, with knowledge of his or her Copyright and Related Rights in the
Work and the meaning and intended legal effect of CC0 on those rights.

1. Copyright and Related Rights. A Work made available under CC0 may be
protected by copyright and related or neighboring rights ("Copyright and
Related Rights"). Copyright and Related Rights include, but are not
limited to, the following:

i. the right to reproduce, adapt, distribute, perform, display,
communicate, and translate a Work;
ii. moral rights retained by the original author(s) and/or performer(s);
iii. publicity and privacy rights pertaining to a person's image or
likeness depicted in a Work;
iv. rights protecting against unfair competition in regards to a Work,
subject to the limitations in paragraph 4(a), below;
v. rights protecting the extraction, dissemination, use and reuse of data
in a Work;
vi. database rights (such as those arising under Directive 96/9/EC of the
European Parliament and of the Council of 11 March 1996 on the legal
protection of databases, and under any national implementation
thereof, including any amended or successor version of such
directive); and
vii. other similar, equivalent or corresponding rights throughout the
world based on applicable law or treaty, and any national
implementations thereof.

2. Waiver. To the greatest extent permitted by, but not in contravention
of, applicable law, Affirmer hereby overtly, fully, permanently,
irrevocably and unconditionally waives, abandons, and surrenders all of
Affirmer's Copyright and Related Rights and associated claims and causes
of action, whether now known or unknown (including existing as well as
future claims and causes of action), in the Work (i) in all territories
worldwide, (ii) for the maximum duration provided by applicable law or
treaty (including future time extensions), (iii) in any current or future
medium and for any number of copies, and (iv) for any purpose whatsoever,
including without limitation commercial, advertising or promotional
purposes (the "Waiver"). Affirmer makes the Waiver for the benefit of each
member of the public at large and to the detriment of Affirmer's heirs and
successors, fully intending that such Waiver shall not be subject to
revocation, rescission, cancellation, termination, or any other legal or
equitable action to disrupt the quiet enjoyment of the Work by the public
as contemplated by Affirmer's express Statement of Purpose.

3. Public License Fallback. Should any part of the Waiver for any reason
be judged legally invalid or ineffective under applicable law, then the
Waiver shall be preserved to the maximum extent permitted taking into
account Affirmer's express Statement of Purpose. In addition, to the
extent the Waiver is so judged Affirmer hereby grants to each affected
person a royalty-free, non transferable, non sublicensable, non exclusive,
irrevocable and unconditional license to exercise Affirmer's Copyright and
Related Rights in the Work (i) in all territories worldwide, (ii) for the
maximum duration provided by applicable law or treaty (including future
time extensions), (iii) in any current or future medium and for any number
of copies, and (iv) for any purpose whatsoever, including without
limitation commercial, advertising or promotional purposes (the
"License"). The License shall be deemed effective as of the date CC0 was
applied by Affirmer to the Work. Should any part of the License for any
reason be judged legally invalid or ineffective under applicable law, such
partial invalidity or ineffectiveness shall not invalidate the remainder
of the License, and in such case Affirmer hereby affirms that he or she
will not (i) exercise any of his or her remaining Copyright and Related
Rights in the Work or (ii) assert any associated claims and causes of
action with respect to the Work, in either case contrary to Affirmer's
express Statement of Purpose.

4. Limitations and Disclaimers.

a. No trademark or patent rights held by Affirmer are waived, abandoned,
surrendered, licensed or otherwise affected by this document.
b. Affirmer offers the Work as-is and makes no representations or
warranties of any kind concerning the Work, express, implied,
statutory or otherwise, including without limitation warranties of
title, merchantability, fitness for a particular purpose, non
infringement, or the absence of latent or other defects, accuracy, or
the present or absence of errors, whether or not discoverable, all to
the greatest extent permissible under applicable law.
c. Affirmer disclaims responsibility for clearing rights of other persons
that may apply to the Work or any use thereof, including without
limitation any person's Copyright and Related Rights in the Work.
Further, Affirmer disclaims responsibility for obtaining any necessary
consents, permissions or other rights required for any use of the
Work.
d. Affirmer understands and acknowledges that Creative Commons is not a
party to this document and has no duty or obligation with respect to
this CC0 or use of the Work.
Loading