diff --git a/Makefile b/Makefile index 58b084c66..3a6e6fc6f 100644 --- a/Makefile +++ b/Makefile @@ -122,7 +122,7 @@ build: # Build fbsql build-fbsql: - $(GO) build -tags='$(BUILD_TAGS)' -ldflags $(LDFLAGS) $(FLAGS) ./cmd/fbsql + CGO_ENABLED=1 $(GO) build -tags='$(BUILD_TAGS)' -ldflags $(LDFLAGS) $(FLAGS) ./cmd/fbsql package: @@ -167,7 +167,7 @@ install-idk: $(MAKE) -C ./idk install install-fbsql: - $(GO) install ./cmd/fbsql + CGO_ENABLED=1 $(GO) install ./cmd/fbsql # Build the lattice assets build-lattice: diff --git a/batch/batcher.go b/batch/batcher.go new file mode 100644 index 000000000..596c51632 --- /dev/null +++ b/batch/batcher.go @@ -0,0 +1,20 @@ +package batch + +import ( + "time" + + "github.com/featurebasedb/featurebase/v3/dax" +) + +// Batcher is an interface implemented by anything which can allocate new +// batches. +type Batcher interface { + NewBatch(cfg Config, tbl *dax.Table, fields []*dax.Field) (RecordBatch, error) +} + +// Config is the configuration options passed to NewBatch for any implementation +// of the Batcher interface. +type Config struct { + Size int + MaxStaleness time.Duration +} diff --git a/cli/batch/inserter.go b/cli/batch/inserter.go new file mode 100644 index 000000000..6401e8523 --- /dev/null +++ b/cli/batch/inserter.go @@ -0,0 +1,8 @@ +package batch + +// Inserter can be implemented by anything which can handle a SQL statement +// representing a write operation. An example is `BULK INSERT`. The Insert() +// method on this interface does not return any results other than an error. +type Inserter interface { + Insert(sql string) error +} diff --git a/cli/batch/sql.go b/cli/batch/sql.go new file mode 100644 index 000000000..e25ef067d --- /dev/null +++ b/cli/batch/sql.go @@ -0,0 +1,215 @@ +package batch + +import ( + "encoding/json" + "fmt" + "strings" + "time" + + fbbatch "github.com/featurebasedb/featurebase/v3/batch" + "github.com/featurebasedb/featurebase/v3/dax" + "github.com/featurebasedb/featurebase/v3/errors" + "github.com/featurebasedb/featurebase/v3/pql" +) + +// Ensure type implements interface. +var _ fbbatch.Batcher = (*sqlBatcher)(nil) + +type sqlBatcher struct { + inserter Inserter + fields []*dax.Field +} + +func NewSQLBatcher(i Inserter, flds []*dax.Field) *sqlBatcher { + return &sqlBatcher{ + inserter: i, + fields: flds, + } +} + +func (b *sqlBatcher) NewBatch(cfg fbbatch.Config, tbl *dax.Table, flds []*dax.Field) (fbbatch.RecordBatch, error) { + fields := flds + if b.fields != nil { + fields = b.fields + } + return &sqlBatch{ + table: tbl, + fields: fields, + size: cfg.Size, + maxStaleness: cfg.MaxStaleness, + ids: make([]interface{}, 0, cfg.Size), + rows: make([][]interface{}, 0, cfg.Size), + inserter: b.inserter, + }, nil +} + +// Ensure type implements interface. +var _ fbbatch.RecordBatch = (*sqlBatch)(nil) + +type sqlBatch struct { + table *dax.Table + fields []*dax.Field + size int + + ids []interface{} + rows [][]interface{} + + // staleTime tracks the time the first record of the batch was inserted + // plus the maxStaleness, in order to raise ErrBatchNowStale if the + // maxStaleness has elapsed + staleTime time.Time + maxStaleness time.Duration + + // inserter handles SQL INSERT statements generated for each batch. + inserter Inserter +} + +func (b *sqlBatch) Add(rec fbbatch.Row) error { + // Clear rec.Values and rec.Clears upon return. + defer func() { + for i := range rec.Values { + rec.Values[i] = nil + } + for k := range rec.Clears { + delete(rec.Clears, k) + } + }() + + if len(b.ids) == cap(b.ids) { + return fbbatch.ErrBatchAlreadyFull + } + if len(rec.Values) != len(b.fields) { + return errors.Errorf("record needs to match up with batch fields, got %d fields and %d record", len(b.fields), len(rec.Values)) + } + + // Append the ID to b.ids. + b.ids = append(b.ids, rec.ID) + + // Convert decimal fields (which come in as int64, along with the scale in + // field) to pql.Decimal. + for i, fld := range b.fields { + switch b.fields[i].Type { + case dax.BaseTypeDecimal: + if val, ok := rec.Values[i].(int64); ok { + rec.Values[i] = pql.NewDecimal(val, fld.Options.Scale) + } + case dax.BaseTypeTimestamp: + if val, ok := rec.Values[i].(int64); ok { + ts := time.Unix(val, 0) + rec.Values[i] = ts.Format(time.RFC3339) + } + } + } + + // Append the record to b.rows. + vals := make([]interface{}, 0, len(rec.Values)) + vals = append(vals, rec.Values...) + b.rows = append(b.rows, vals) + + // Check for batch full or stale. + if len(b.ids) == cap(b.ids) { + return fbbatch.ErrBatchNowFull + } + if b.maxStaleness != time.Duration(0) { // set maxStaleness to 0 to disable staleness checking + if len(b.ids) == 1 { + b.staleTime = time.Now().Add(b.maxStaleness) + } else if time.Now().After(b.staleTime) { + return fbbatch.ErrBatchNowStale + } + } + return nil +} + +func (b *sqlBatch) Import() error { + if len(b.rows) == 0 { + return nil + } + + // Construct the BULK INSERT statement based on the table and fields. + sql, err := buildBulkInsert(b.table, b.fields, b.ids, b.rows) + if err != nil { + return errors.Wrap(err, "building bulk insert statement") + } + + // Reset batch data. + b.reset() + + // Submit the SQL statement. + return b.inserter.Insert(sql) +} + +func (b *sqlBatch) reset() { + b.ids = b.ids[:0] + b.rows = b.rows[:0] +} + +func (b *sqlBatch) Len() int { + return len(b.rows) +} + +func (b *sqlBatch) Flush() error { + return nil +} + +func buildBulkInsert(tbl *dax.Table, fields []*dax.Field, ids []interface{}, rows [][]interface{}) (string, error) { + // Validation. + if tbl.Name == "" { + return "", errors.New(errors.ErrUncoded, "table name is required") + } else if len(fields) == 0 { + return "", errors.New(errors.ErrUncoded, "at least one field is required") + } + + var sb strings.Builder + + sb.WriteString(`BULK INSERT INTO `) + sb.WriteString(string(tbl.Name)) + sb.WriteString(` (_id,`) + + flds := make([]string, 0, len(fields)) + maps := make([]string, 0, len(fields)) + for i := range fields { + flds = append(flds, string(fields[i].Name)) + maps = append(maps, fmt.Sprintf("'$.col_%d' %s", i, fields[i].Definition())) + } + // Fields + sb.WriteString(strings.Join(flds, ",")) + + // MAP + keyType := dax.BaseTypeID + if tbl.StringKeys() { + keyType = dax.BaseTypeString + } + sb.WriteString(`) MAP ('$._id' `) + sb.WriteString(keyType) + sb.WriteString(`,`) + sb.WriteString(strings.Join(maps, ",")) + sb.WriteString(`) FROM x'`) + + // Row values. + + // m is a map representing a single row to be marshalled and added to the + // bulk insert as one line in the NDJSON payload. We re-use the map for each + // row. + m := make(map[string]interface{}) + for i := range rows { + // Write the ID value. + m[string(dax.PrimaryKeyFieldName)] = ids[i] + // Write the rest of the data values. + for col := range rows[i] { + m[fmt.Sprintf("col_%d", col)] = rows[i][col] + } + + // Marshal the map to json and add to the sql statement. + if j, err := json.Marshal(m); err != nil { + return "", errors.Wrap(err, "marshalling row to json") + } else { + sb.Write(j) + sb.WriteString("\n") + } + } + + // WITH + sb.WriteString(fmt.Sprintf(`' WITH BATCHSIZE %d FORMAT 'NDJSON' INPUT 'STREAM'`, len(rows))) + + return sb.String(), nil +} diff --git a/cli/batch/sql_test.go b/cli/batch/sql_test.go new file mode 100644 index 000000000..86aace469 --- /dev/null +++ b/cli/batch/sql_test.go @@ -0,0 +1,47 @@ +package batch + +import ( + "testing" + + "github.com/featurebasedb/featurebase/v3/dax" + "github.com/stretchr/testify/assert" +) + +func TestBatchSQL(t *testing.T) { + tbl := &dax.Table{ + Name: "foo", + } + fields := []*dax.Field{ + { + Name: "name", + Type: dax.BaseTypeString, + }, + { + Name: "age", + Type: dax.BaseTypeInt, + }, + } + ids := []interface{}{ + 0, 1, 2, + } + rows := [][]interface{}{ + { + []interface{}{"Alice", int64(11)}, + }, + { + []interface{}{"Bob", int64(22)}, + }, + { + []interface{}{"Carl,Comma", int64(33)}, + }, + } + + s, err := buildBulkInsert(tbl, fields, ids, rows) + assert.NoError(t, err) + + exp := `BULK INSERT INTO foo (_id,name,age) MAP (0 int,1 string,2 int) FROM x'0,Alice,11 +1,Bob,22 +2,"Carl,Comma",33 +' WITH BATCHSIZE 3 FORMAT 'CSV' INPUT 'STREAM'` + assert.Equal(t, exp, s) +} diff --git a/cli/cli.go b/cli/cli.go index e684961ed..18e412c90 100644 --- a/cli/cli.go +++ b/cli/cli.go @@ -13,9 +13,12 @@ import ( "github.com/chzyer/readline" featurebase "github.com/featurebasedb/featurebase/v3" + "github.com/featurebasedb/featurebase/v3/cli/batch" "github.com/featurebasedb/featurebase/v3/cli/fbcloud" + "github.com/featurebasedb/featurebase/v3/cli/kafka" "github.com/featurebasedb/featurebase/v3/errors" "github.com/featurebasedb/featurebase/v3/logger" + "github.com/spf13/viper" ) const ( @@ -38,6 +41,7 @@ Type "\q" to quit. // Ensure type implments interfaces. var _ printer = (*Command)(nil) +var _ batch.Inserter = (*Command)(nil) type Command struct { host string @@ -85,6 +89,10 @@ type Command struct { // `-c` flag in the command line. nonInteractiveMode bool + // kafkaRunner will be non-nil if the CLI has been configured to consume + // from kafka. In that case, Command will run in non-interactive mode. + kafkaRunner *kafka.Runner + // quit gets closed when Run should stop listening for input. quit chan struct{} } @@ -129,7 +137,21 @@ func NewCommand(logdest logger.Logger) *Command { // Run is the main entry-point to the CLI. func (cmd *Command) Run(ctx context.Context) error { - cmd.setupConfig() + if err := cmd.run(ctx); err != nil { + cmd.Errorf(err.Error() + "\n") + return err + } + return nil +} + +// run is effectively wrapped by the Run() method, but it's split out this way +// so that run() can simply return errors, rather than worrying about how errors +// should be printed; printing errors returned by run() is left up to the Run() +// method. +func (cmd *Command) run(ctx context.Context) error { + if err := cmd.setupConfig(); err != nil { + return errors.Wrap(err, "setting up config") + } // Check to see if Command needs to run in non-interactive mode. if len(cmd.Commands) > 0 || len(cmd.Files) > 0 { @@ -140,24 +162,36 @@ func (cmd *Command) Run(ctx context.Context) error { } if err := cmd.connectToDatabase(cmd.database); err != nil { cmd.Errorf(errors.Wrap(err, "connecting to database").Error() + "\n") + // We intentionally do not return err here. } // Run Commands. for _, line := range cmd.Commands { if err := cmd.handleLine(line); err != nil { - cmd.Errorf(err.Error()) - return nil + return errors.Wrapf(err, "handling line: %s", line) } } // Run Files. for _, fname := range cmd.Files { if _, err := executeFile(cmd, fname); err != nil { - cmd.Errorf(err.Error()) - return nil + return errors.Wrapf(err, "executing file: %s", fname) } } + return nil + } else if cmd.kafkaRunner != nil { + if err := cmd.setupClient(); err != nil { + return errors.Wrap(err, "setting up client") + } + if err := cmd.connectToDatabase(cmd.database); err != nil { + cmd.Errorf(errors.Wrap(err, "connecting to database").Error() + "\n") + // We intentionally do not return err here. + } + + if err := cmd.kafkaRunner.Main.Run(); err != nil { + return errors.Wrap(err, "running kafka") + } return nil } @@ -170,6 +204,7 @@ func (cmd *Command) Run(ctx context.Context) error { cmd.printConnInfo() if err := cmd.connectToDatabase(cmd.database); err != nil { cmd.Errorf(errors.Wrap(err, "connecting to database").Error() + "\n") + // We intentionally do not return err here. } rl, err := readline.NewEx(&readline.Config{ @@ -308,9 +343,9 @@ func (cmd *Command) close() error { // setupConfig sets up private struct members based on values provided via the // configuration flags. -func (cmd *Command) setupConfig() { +func (cmd *Command) setupConfig() error { if cmd.Config == nil { - return + return nil } cmd.host = cmd.Config.Host @@ -320,6 +355,39 @@ func (cmd *Command) setupConfig() { cmd.database = cmd.Config.Database cmd.historyPath = cmd.Config.HistoryPath + + // Kafka setup. + if cmd.Config.KafkaConfig != "" { + // read the kafka config file + v := viper.New() + v.SetConfigFile(cmd.Config.KafkaConfig) + v.SetConfigType("toml") + err := v.ReadInConfig() + if err != nil { + return fmt.Errorf("error reading configuration file '%s': %v", cmd.Config.KafkaConfig, err) + } + + cfg := kafka.Config{} + if err := v.Unmarshal(&cfg); err != nil { + return errors.Wrap(err, "unmarshalling config") + } + + if err := kafka.ValidateConfig(cfg); err != nil { + return errors.Wrap(err, "validating config") + } + + cleanCfg, err := kafka.ConvertConfig(cfg) + if err != nil { + return errors.Wrap(err, "cleaning config") + } + + cmd.kafkaRunner = kafka.NewRunner( + cleanCfg, + batch.NewSQLBatcher(cmd, kafka.ConfigToFields(cfg)), + ) + } + + return nil } func (cmd *Command) executeAndWriteQuery(qry query) error { @@ -462,7 +530,7 @@ func (cmd *Command) connectionMessage() string { if cmd.databaseName == "" { return "You are not connected to a database.\n" } - return fmt.Sprintf("You are now connected to database \"%s\" (%s) as user \"???\".\n", cmd.databaseName, cmd.databaseID) + return fmt.Sprintf("You are now connected to database \"%s\" (%s).\n", cmd.databaseName, cmd.databaseID) } func (cmd *Command) setupClient() error { @@ -695,3 +763,11 @@ func (cmd *Command) handleLineAsQueryParts(line string) error { } return nil } + +func (cmd *Command) Insert(sql string) error { + wqr, err := cmd.executeQuery(newRawQuery(sql)) + if wqr.Error != "" { + return errors.Errorf(wqr.Error) + } + return err +} diff --git a/cli/config.go b/cli/config.go index 3953388c1..6d9722871 100644 --- a/cli/config.go +++ b/cli/config.go @@ -11,6 +11,9 @@ type Config struct { // CloudAuth CloudAuth CloudAuthConfig `json:"cloud-auth"` + // Kafka + KafkaConfig string `json:"kafka-config"` + HistoryPath string `json:"history-path"` } diff --git a/cli/kafka/config.go b/cli/kafka/config.go new file mode 100644 index 000000000..8b46ccd01 --- /dev/null +++ b/cli/kafka/config.go @@ -0,0 +1,172 @@ +package kafka + +import ( + "fmt" + "time" + + "github.com/featurebasedb/featurebase/v3/dax" + "github.com/featurebasedb/featurebase/v3/idk" + "github.com/pkg/errors" +) + +type Config struct { + Hosts []string `mapstructure:"hosts" help: "Kafka hosts."` + Group string `mapstructure:"group" help:"Kafka group."` + Topics []string `mapstructure:"topics" help:"Kafka topics to read from."` + + BatchSize int `mapstructure:"batch-size" help:"Batch size."` + BatchMaxStaleness time.Duration `mapstructure:"batch-max-staleness" help:"Maximum length of time that the oldest record in a batch can exist before flushing the batch. Note that this can potentially stack with timeouts waiting for the source."` + Timeout time.Duration `mapstructure:"timeout" help:"Time to wait for more records from Kafka before flushing a batch. 0 to disable."` + + Table string `mapstructure:"table" help:"Destination table name."` + Fields []Field `mapstructure:"fields"` +} + +type ConfigForIDK struct { + Hosts []string + Group string + Topics []string + + BatchSize int + BatchMaxStaleness time.Duration + Timeout time.Duration + + Table string + IDField string + Fields []idk.RawField +} + +type Field struct { + Name string `mapstructure:"name"` + Type string `mapstructure:"type"` + SourcePath []string `mapstructure:"source-path"` + PrimaryKey bool `mapstructure:"primary-key"` + + Options FieldOptions `mapstructure:"options"` +} + +type FieldOptions struct { + Scale int64 `mapstructure:"scale"` +} + +// ValidateConfig validates the config is usable. +func ValidateConfig(c Config) error { + if c.Table == "" { + return errors.Errorf("table is required") + } else if len(c.Topics) == 0 { + return errors.Errorf("at least one topic is required") + } else if len(c.Fields) < 2 { + return errors.Errorf("at least two fields are required (one should be a primary key)") + } else { + var found int + for i := range c.Fields { + if c.Fields[i].PrimaryKey { + found++ + } + } + if found != 1 { + return errors.Errorf("exactly one primary key field is required") + } + } + return nil +} + +// ConvertConfig converts a Config to one that suitable for IDK. +func ConvertConfig(c Config) (ConfigForIDK, error) { + // Set a default kafka host in case one isn't provided. + hosts := []string{"localhost:9092"} + if len(c.Hosts) > 0 { + hosts = c.Hosts + } + + // Copy all the shared members from Config to ConfigForIDK. + out := ConfigForIDK{ + Hosts: hosts, + Group: c.Group, + Topics: c.Topics, + BatchSize: c.BatchSize, + BatchMaxStaleness: c.BatchMaxStaleness, + Timeout: c.Timeout, + Table: c.Table, + } + + if len(c.Fields) == 0 { + return out, errors.New("fields cannot be empty") + } + + // rawFields wil be the same as c.Fields, but possibly enhanced. + rawFields := make([]idk.RawField, 0, len(c.Fields)) + + var foundPK bool + for _, fld := range c.Fields { + if fld.PrimaryKey { + out.IDField = fld.Name + foundPK = true + } + + typ, err := dax.BaseTypeFromString(fld.Type) + if err != nil { + return out, errors.Wrap(err, "getting base type") + } + + rawFld := idk.RawField{ + Name: fld.Name, + Type: string(typ), + Path: fld.SourcePath, + } + // If a SourcePath wasn't provided, default to using the field name. + if len(rawFld.Path) == 0 { + rawFld.Path = []string{fld.Name} + } + + switch typ { + case dax.BaseTypeInt: + // We don't have to handle min/max because we don't create the table. + case dax.BaseTypeDecimal: + rawFld.Config = []byte(fmt.Sprintf(`{"scale":%d}`, fld.Options.Scale)) + case dax.BaseTypeID: + rawFld.Config = []byte("{\"mutex\":true}") + case dax.BaseTypeIDSet: + rawFld.Type = "ids" + case dax.BaseTypeString: + rawFld.Config = []byte("{\"mutex\":true}") + case dax.BaseTypeStringSet: + rawFld.Type = "strings" + case dax.BaseTypeTimestamp: + // No timestamp options are handled for now. + } + + rawFields = append(rawFields, rawFld) + } + if !foundPK { + return out, errors.New("primary-key not found in fields") + } + + out.Fields = rawFields + + return out, nil +} + +// ConfigToFields returns a list of *dax.Field based on the IDField and Fields +// in the Config. +func ConfigToFields(c Config) []*dax.Field { + // We don't know if a primary key will be found, so we can't set the + // capacity to `len(c.Fields)-1`. + out := make([]*dax.Field, 0, len(c.Fields)) + + for _, fld := range c.Fields { + if fld.PrimaryKey { + continue + } + dfld := &dax.Field{ + Name: dax.FieldName(fld.Name), + Type: dax.BaseType(fld.Type), + Options: dax.FieldOptions{ + Scale: fld.Options.Scale, + }, + } + out = append(out, dfld) + } + + return out +} diff --git a/cli/kafka/runner.go b/cli/kafka/runner.go new file mode 100644 index 000000000..739b2bd2b --- /dev/null +++ b/cli/kafka/runner.go @@ -0,0 +1,61 @@ +package kafka + +import ( + "time" + + fbbatch "github.com/featurebasedb/featurebase/v3/batch" + "github.com/featurebasedb/featurebase/v3/errors" + "github.com/featurebasedb/featurebase/v3/idk" + "github.com/featurebasedb/featurebase/v3/idk/kafka_static" +) + +type Runner struct { + idk.Main `flag:"!embed"` + KafkaHosts []string `help:"Comma separated list of host:port pairs for Kafka."` + Group string `help:"Kafka group."` + Topics []string `help:"Kafka topics to read from."` + Timeout time.Duration `help:"Time to wait for more records from Kafka before flushing a batch. 0 to disable."` + Header []idk.RawField `help:"Header configuration."` +} + +func NewRunner(cfg ConfigForIDK, batcher fbbatch.Batcher) *Runner { + idkMain := idk.NewMain() + idkMain.IDField = cfg.IDField + idkMain.Index = cfg.Table + idkMain.Batcher = batcher + idkMain.BatchSize = cfg.BatchSize + idkMain.BatchMaxStaleness = cfg.BatchMaxStaleness + idkMain.Basic() + + kr := &Runner{ + Main: *idkMain, + KafkaHosts: cfg.Hosts, + Group: cfg.Group, + Topics: cfg.Topics, + Header: cfg.Fields, + Timeout: cfg.Timeout, + } + kr.OffsetMode = true + kr.Main.Namespace = "cli_kafka_runner" + kr.Main.Pprof = "" // don't initialize pprof until we actually use it in tests + kr.NewSource = func() (idk.Source, error) { + source := kafka_static.NewSource() + source.Hosts = kr.KafkaHosts + source.Group = kr.Group + source.Topics = kr.Topics + source.Log = kr.Main.Log() + // source.TLS = m.KafkaTLS + source.Timeout = kr.Timeout + // source.SkipOld = m.SkipOld + source.HeaderFields = kr.Header + // source.S3Region = m.S3Region + // source.AllowMissingFields = m.AllowMissingFields + + err := source.Open() + if err != nil { + return nil, errors.Wrap(err, "opening source") + } + return source, nil + } + return kr +} diff --git a/cli/parts.go b/cli/parts.go index a936ebc59..fd704cae6 100644 --- a/cli/parts.go +++ b/cli/parts.go @@ -39,6 +39,12 @@ type queryPart interface { Reader() io.Reader } +func newRawQuery(s string) query { + return []queryPart{ + newPartRaw(s), + } +} + // //////////////////////////////////////////////////////////////////////////// // raw // //////////////////////////////////////////////////////////////////////////// diff --git a/cmd/auth_token.go b/cmd/auth_token.go index 7921c1f91..1fb52a20b 100644 --- a/cmd/auth_token.go +++ b/cmd/auth_token.go @@ -16,7 +16,7 @@ func newAuthTokenCommand(logdest logger.Logger) *cobra.Command { Long: ` Retrieves an auth-token for use in authenticating with FeatureBase from the configured identity provider. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/backup.go b/cmd/backup.go index cf962fc42..d5234c5eb 100644 --- a/cmd/backup.go +++ b/cmd/backup.go @@ -16,7 +16,7 @@ func newBackupCommand(logdest logger.Logger) *cobra.Command { Long: ` Backs up a FeatureBase server to a local, tar-formatted snapshot file. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/backup_tar.go b/cmd/backup_tar.go index 641f4ee51..3deaf8039 100644 --- a/cmd/backup_tar.go +++ b/cmd/backup_tar.go @@ -16,7 +16,7 @@ func newBackupTarCommand(logdest io.Writer) *cobra.Command { Long: ` Backs up a FeatureBase server to a local, tar-formatted snapshot file. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/chksum.go b/cmd/chksum.go index 4160114f4..8edaaf458 100644 --- a/cmd/chksum.go +++ b/cmd/chksum.go @@ -19,7 +19,7 @@ func newChkSumCommand(logdest logger.Logger) *cobra.Command { Generates a digital signature of all the data associated with a provided FeatureBase server WARNING: could be slow if high cardinality fields exist `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/cli.go b/cmd/cli.go deleted file mode 100644 index 09e7428d4..000000000 --- a/cmd/cli.go +++ /dev/null @@ -1,35 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -package cmd - -import ( - "io" - - "github.com/featurebasedb/featurebase/v3/cli" - "github.com/featurebasedb/featurebase/v3/ctl" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/spf13/cobra" - "github.com/spf13/viper" -) - -var cliCmd *cli.Command - -// NewCLICommand runs the FeatureBase CLI subcommand. -func NewCLICommand(stderr io.Writer) *cobra.Command { - logdest := logger.NewStandardLogger(stderr) - cliCmd = cli.NewCommand(logdest) - cobraCmd := &cobra.Command{ - Use: "fbsql", - Short: "Query FeatureBase with SQL from the command line", - Long: ``, - RunE: usageErrorWrapper(cliCmd), - PersistentPreRunE: func(cmd *cobra.Command, args []string) error { - v := viper.New() - return setAllConfig(v, cmd.Flags(), "FBSQL") - }, - SilenceErrors: true, - } - - // Attach flags to the command. - ctl.BuildCLIFlags(cobraCmd, cliCmd) - return cobraCmd -} diff --git a/cmd/dataframe-csv-loader.go b/cmd/dataframe-csv-loader.go index fd0679b1d..39b77b0bb 100644 --- a/cmd/dataframe-csv-loader.go +++ b/cmd/dataframe-csv-loader.go @@ -15,7 +15,7 @@ func newDataframeCsvLoaderCommand(logdest logger.Logger) *cobra.Command { Short: "load dataframe integer and floating point values into featurebase", Long: ` `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := loaderCmd.Flags() flags.StringVar(&cmd.Path, "csv", "", "path to csv input file") diff --git a/cmd/export.go b/cmd/export.go index 88c817fd5..dbb0ec77e 100644 --- a/cmd/export.go +++ b/cmd/export.go @@ -26,7 +26,7 @@ The format of the CSV file is: The file does not contain any headers. `, - RunE: usageErrorWrapper(Exporter), + RunE: UsageErrorWrapper(Exporter), } flags := exportCmd.Flags() diff --git a/cmd/fbsql/main.go b/cmd/fbsql/main.go index d498d1703..07655bcd5 100644 --- a/cmd/fbsql/main.go +++ b/cmd/fbsql/main.go @@ -2,12 +2,63 @@ package main import ( + "io" "os" + "github.com/featurebasedb/featurebase/v3/cli" "github.com/featurebasedb/featurebase/v3/cmd" + "github.com/featurebasedb/featurebase/v3/logger" + "github.com/spf13/cobra" + "github.com/spf13/viper" ) func main() { - command := cmd.NewCLICommand(os.Stderr) + command := newCLICommand(os.Stderr) command.Execute() } + +// newCLICommand runs the FeatureBase CLI subcommand. +func newCLICommand(stderr io.Writer) *cobra.Command { + logdest := logger.NewStandardLogger(stderr) + cliCmd := cli.NewCommand(logdest) + cobraCmd := &cobra.Command{ + Use: "fbsql", + Short: "Query FeatureBase with SQL from the command line", + Long: ``, + RunE: cmd.UsageErrorWrapper(cliCmd), + PersistentPreRunE: func(cobraCmd *cobra.Command, args []string) error { + v := viper.New() + return cmd.SetAllConfig(v, cobraCmd.Flags(), "FBSQL") + }, + SilenceErrors: true, + } + + // Attach flags to the command. + buildFlags(cobraCmd, cliCmd) + return cobraCmd +} + +// buildFlags attaches a set of flags to the command for a cli instance. +func buildFlags(cmd *cobra.Command, cliCmd *cli.Command) { + flags := cmd.Flags() + + // Base struct flags. + flags.StringSliceVarP(&cliCmd.Commands, "command", "c", cliCmd.Commands, "Command to run in non-interactive mode. Provide multiple flags to execute more than one command. All `--command` flags run before all `--file` flags.") + flags.StringSliceVarP(&cliCmd.Files, "file", "f", cliCmd.Files, "File to run in non-interactive mode. Provide multiple flags to execute more than one file. All `--command` flags run before all `--file` flags.") + + // Config flags. + flags.StringVarP(&cliCmd.Config.Host, "host", "", cliCmd.Config.Host, "hostname of FeatureBase.") + flags.StringVarP(&cliCmd.Config.Port, "port", "", cliCmd.Config.Port, "port of FeatureBase.") + flags.StringVar(&cliCmd.Config.HistoryPath, "history-path", cliCmd.Config.HistoryPath, "path for history files.") + flags.StringVar(&cliCmd.Config.OrganizationID, "org-id", cliCmd.Config.OrganizationID, "OrganizationID.") + flags.StringVarP(&cliCmd.Config.Database, "dbname", "d", cliCmd.Config.Database, "Name of the database to connect to.") + + flags.StringVar(&cliCmd.Config.CloudAuth.ClientID, "client-id", cliCmd.Config.CloudAuth.ClientID, "Cognito Client ID for FeatureBase Cloud access.") + flags.StringVar(&cliCmd.Config.CloudAuth.Region, "region", cliCmd.Config.CloudAuth.Region, "Cloud region for FeatureBase Cloud access (e.g. us-east-2).") + flags.StringVar(&cliCmd.Config.CloudAuth.Email, "email", cliCmd.Config.CloudAuth.Email, "Email address for FeatureBase Cloud access.") + flags.StringVar(&cliCmd.Config.CloudAuth.Password, "password", cliCmd.Config.CloudAuth.Password, "Password for FeatureBase Cloud access.") + + flags.StringVar(&cliCmd.Config.KafkaConfig, "kafka-config", cliCmd.Config.KafkaConfig, "Kafka configuration file to read from.") + + flags.String("config", "", "Configuration file to read from.") +} diff --git a/cmd/generate_config.go b/cmd/generate_config.go index 012da9fd3..882b51bbf 100644 --- a/cmd/generate_config.go +++ b/cmd/generate_config.go @@ -18,7 +18,7 @@ func newGenerateConfigCommand(logdest logger.Logger) *cobra.Command { Short: "Print the default configuration.", Long: `generate-config prints the default configuration to stdout `, - RunE: usageErrorWrapper(generateConf), + RunE: UsageErrorWrapper(generateConf), } return confCmd diff --git a/cmd/keygen.go b/cmd/keygen.go index 4dc25df38..ae8c2d60f 100644 --- a/cmd/keygen.go +++ b/cmd/keygen.go @@ -16,7 +16,7 @@ func newKeygenCommand(logdest logger.Logger) *cobra.Command { Long: ` Generate secret key to configure FeatureBase for Authentication. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/parquet-info.go b/cmd/parquet-info.go index 746b1ec21..12c496e79 100644 --- a/cmd/parquet-info.go +++ b/cmd/parquet-info.go @@ -27,7 +27,7 @@ Displays schema and sample data from the specified file c.Path = args[0] return nil }, - RunE: usageErrorWrapper(c), + RunE: UsageErrorWrapper(c), } return cmd } diff --git a/cmd/presort.go b/cmd/presort.go index 1bd4aba21..23c48965b 100644 --- a/cmd/presort.go +++ b/cmd/presort.go @@ -15,7 +15,7 @@ func newPreSortCommand(logdest logger.Logger) *cobra.Command { Long: ` Takes all input files and writes PartitionN numbered files to a directory, where each file contains only records that will go into the partition it is named for. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := ccmd.Flags() diff --git a/cmd/rbf.go b/cmd/rbf.go index 500d54dbe..2ca84be87 100644 --- a/cmd/rbf.go +++ b/cmd/rbf.go @@ -44,7 +44,7 @@ Executes a consistency check on an RBF data directory. c.Path = args[0] return nil }, - RunE: usageErrorWrapper(c), + RunE: UsageErrorWrapper(c), } return cmd } @@ -76,7 +76,7 @@ Dumps the raw hex data for one or more RBF pages. return nil }, - RunE: usageErrorWrapper(c), + RunE: UsageErrorWrapper(c), } return cmd } @@ -98,7 +98,7 @@ Prints a line for every page in the database with its type/status. c.Path = args[0] return nil }, - RunE: usageErrorWrapper(c), + RunE: UsageErrorWrapper(c), } flags := cmd.Flags() @@ -133,7 +133,7 @@ Prints the header & cell data for one or more pages. return nil }, - RunE: usageErrorWrapper(c), + RunE: UsageErrorWrapper(c), } return cmd } diff --git a/cmd/restore.go b/cmd/restore.go index f94b9de71..931b117e2 100644 --- a/cmd/restore.go +++ b/cmd/restore.go @@ -16,7 +16,7 @@ func newRestoreCommand(logdest logger.Logger) *cobra.Command { Long: ` The Restore command will take a backup archive and restore it to a new, clean cluster. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := restoreCmd.Flags() flags.StringVarP(&cmd.Path, "source", "s", "", "backup file; specify '-' to restore from stdin tar stream") diff --git a/cmd/restore_tar.go b/cmd/restore_tar.go index 9a1cdc3fc..32b323926 100644 --- a/cmd/restore_tar.go +++ b/cmd/restore_tar.go @@ -15,7 +15,7 @@ func newRestoreTarCommand(logdest logger.Logger) *cobra.Command { Long: ` The Restore command will take a tar-formatted backup archive and restore it to a new, clean cluster. `, - RunE: usageErrorWrapper(cmd), + RunE: UsageErrorWrapper(cmd), } flags := restoreCmd.Flags() flags.StringVarP(&cmd.Path, "source", "s", "", "backup file; specify '-' to restore from stdin tar stream") diff --git a/cmd/root.go b/cmd/root.go index 9a67365fc..d05956e9f 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -25,11 +25,11 @@ type runner interface { Run(context.Context) error } -// usageErrorWrapper takes a thing with a Run(context) error, and produces +// UsageErrorWrapper takes a thing with a Run(context) error, and produces // a func(*cobra.Command, []string) error from it which will run that // command, and then set Cobra's SilenceUsage flag unless the returned // error errors.Is() a ctl.UsageError. -func usageErrorWrapper(inner runner) func(*cobra.Command, []string) error { +func UsageErrorWrapper(inner runner) func(*cobra.Command, []string) error { return func(c *cobra.Command, args []string) error { return considerUsageError(c, inner.Run(context.Background())) } @@ -68,7 +68,7 @@ at https://docs.featurebase.com/. case "dax": v.Set("future.rename", true) // always use FEATUREBASE env for dax } - if err := setAllConfig(v, cmd.Flags(), ""); err != nil { + if err := SetAllConfig(v, cmd.Flags(), ""); err != nil { return err } @@ -114,17 +114,17 @@ at https://docs.featurebase.com/. return rc } -// setAllConfig takes a FlagSet to be the definition of all configuration +// SetAllConfig takes a FlagSet to be the definition of all configuration // options, as well as their defaults. It then reads from the command line, the // environment, and a config file (if specified), and applies the configuration // in that priority order. Since each flag in the set contains a pointer to -// where its value should be stored, setAllConfig can directly modify the value +// where its value should be stored, SetAllConfig can directly modify the value // of each config variable. // -// setAllConfig looks for environment variables which are capitalized versions +// SetAllConfig looks for environment variables which are capitalized versions // of the flag names with dashes replaced by underscores, and prefixed with // envPrefix plus an underscore. -func setAllConfig(v *viper.Viper, flags *pflag.FlagSet, envPrefix string) error { // nolint: unparam +func SetAllConfig(v *viper.Viper, flags *pflag.FlagSet, envPrefix string) error { // nolint: unparam // add cmd line flag def to viper err := v.BindPFlags(flags) if err != nil { diff --git a/ctl/cli.go b/ctl/cli.go deleted file mode 100644 index 564f708a5..000000000 --- a/ctl/cli.go +++ /dev/null @@ -1,39 +0,0 @@ -package ctl - -import ( - "github.com/featurebasedb/featurebase/v3/cli" - "github.com/spf13/cobra" - "github.com/spf13/pflag" -) - -// BuildCLIFlags attaches a set of flags to the command for a cli instance. -func BuildCLIFlags(cmd *cobra.Command, cliCmd *cli.Command) { - flags := cmd.Flags() - - // Base struct flags. - flags.StringSliceVarP(&cliCmd.Commands, "command", "c", cliCmd.Commands, "Command to run in non-interactive mode. Provide multiple flags to execute more than one command. All `--command` flags run before all `--file` flags.") - flags.StringSliceVarP(&cliCmd.Files, "file", "f", cliCmd.Files, "File to run in non-interactive mode. Provide multiple flags to execute more than one file. All `--command` flags run before all `--file` flags.") - - // Config flags. - flags.AddFlagSet(cliConfigFlagSet(cliCmd.Config)) -} - -// cliConfigFlagSet returns a pflag.FlagSet for the CLI Config struct. -func cliConfigFlagSet(cfg *cli.Config) *pflag.FlagSet { - flags := pflag.NewFlagSet("cli", pflag.ExitOnError) - - flags.StringVarP(&cfg.Host, "host", "", cfg.Host, "hostname of FeatureBase.") - flags.StringVarP(&cfg.Port, "port", "", cfg.Port, "port of FeatureBase.") - flags.StringVar(&cfg.HistoryPath, "history-path", cfg.HistoryPath, "path for history files.") - flags.StringVar(&cfg.OrganizationID, "org-id", cfg.OrganizationID, "OrganizationID.") - flags.StringVar(&cfg.Database, "db", cfg.Database, "Name of the database to connect to.") - - flags.StringVar(&cfg.CloudAuth.ClientID, "client-id", cfg.CloudAuth.ClientID, "Cognito Client ID for FeatureBase Cloud access.") - flags.StringVar(&cfg.CloudAuth.Region, "region", cfg.CloudAuth.Region, "Cloud region for FeatureBase Cloud access (e.g. us-east-2).") - flags.StringVar(&cfg.CloudAuth.Email, "email", cfg.CloudAuth.Email, "Email address for FeatureBase Cloud access.") - flags.StringVar(&cfg.CloudAuth.Password, "password", cfg.CloudAuth.Password, "Password for FeatureBase Cloud access.") - - flags.String("config", "", "Configuration file to read from.") - - return flags -} diff --git a/dax/table.go b/dax/table.go index 37bcc0179..75057468c 100644 --- a/dax/table.go +++ b/dax/table.go @@ -707,6 +707,17 @@ func (f *Field) String() string { return string(f.Name) } +// Definition returns the field name along with its parenthetical (when +// applicable). +func (f *Field) Definition() string { + switch f.Type { + case BaseTypeDecimal: + return fmt.Sprintf("%s(%d)", f.Type, f.Options.Scale) + default: + return string(f.Type) + } +} + // StringKeys returns true if the field uses string keys. func (f *Field) StringKeys() bool { switch f.Type { diff --git a/idk/header.go b/idk/header.go index cb7a2166d..20912a0ff 100644 --- a/idk/header.go +++ b/idk/header.go @@ -506,13 +506,18 @@ func (t PathTable) FlatMap() map[string]int { return m } +// RawField is used in cases where header fields are configured as json, +// typically read from a file. But this type is also used by the kafka runner in +// fbsql. +type RawField struct { + Name string `json:"name"` + Path []string `json:"path"` + Type string `json:"type"` + Config json.RawMessage +} + func ParseHeader(raw []byte) ([]Field, PathTable, error) { - var rawSchema []struct { - Name string `json:"name"` - Path []string `json:"path"` - Type string `json:"type"` - Config json.RawMessage - } + var rawSchema []RawField err := json.Unmarshal(raw, &rawSchema) if err != nil { return nil, nil, errors.Wrap(err, "parsing schema") diff --git a/idk/ingest.go b/idk/ingest.go index e760c2588..385ff2889 100644 --- a/idk/ingest.go +++ b/idk/ingest.go @@ -124,6 +124,12 @@ type Main struct { NewImporterFn func() pilosacore.Importer `flag:"-"` + Batcher pilosabatch.Batcher + + // basic, when true, will only set up the things required to run a basic + // ingester. For example, it does not set up the pilosa client. + basic bool + SchemaManager SchemaManager `flag:"-"` Qtbl *dax.QualifiedTable `flag:"-"` @@ -241,6 +247,12 @@ func (m *Main) Rename() { } } +// Basic sets up Main with basic functionality, excluding those things which are +// not required for some implementations (such as the kafka runner in fbsql). +func (m *Main) Basic() { + m.basic = true +} + func (m *Main) Run() (err error) { onFinishRun, err := m.Setup() if err != nil { @@ -674,6 +686,100 @@ initialFetch: } func (m *Main) Setup() (onFinishRun func(), err error) { + if m.basic { + return m.basicSetup() + } + return m.setup() +} + +// basicSetup contains a lot of the same functionality as setup(), but it +// exludes anything which involves interacting with a "destination" featurebase +// installation. A basic setup is useful for something which wants to use the +// ingest loop and its batching logic, but doesn't want to send results directly +// to a featurebase installation. An example of this would be the CLI (i.e. +// fbsql), which generates BULK INSERT statements and sends those to a /sql +// endpoint. +func (m *Main) basicSetup() (onFinishRun func(), err error) { + if err := m.validate(); err != nil { + return nil, errors.Wrap(err, "validating configuration") + } + + // setup logging + var f *logger.FileWriter + var logOut io.Writer = os.Stderr + if m.LogPath != "" { + f, err = logger.NewFileWriter(m.LogPath) + if err != nil { + return nil, errors.Wrap(err, "opening log file") + } + logOut = f + } + if m.Verbose { + m.log = logger.NewVerboseLogger(logOut) + } else { + m.log = logger.NewStandardLogger(logOut) + } + + if m.TrackProgress { + m.progress = &ProgressTracker{} + } + + // Set up progress tracking. + if m.progress != nil { + startTime := time.Now() + var wg sync.WaitGroup + defer wg.Wait() + doneCh := make(chan struct{}) + defer func() { close(doneCh) }() + wg.Add(1) + go func() { + defer wg.Done() + + // Set up a timer to check progress every 10 seconds. + tick := time.NewTicker(10 * time.Second) + defer tick.Stop() + + prev := uint64(0) + stalled := true + for { + progress := m.progress.Check() + switch { + case progress != prev: + // Forward progress continues. + m.log.Printf("sourced %d records (%.2f records/minute)", progress, float64(progress)/time.Since(startTime).Minutes()) + stalled = false + case stalled: + // We already told the user that it is stalled. + default: + // This is the start of a stall. + // No records have been sourced in the past 5 seconds. + m.log.Printf("record sourcing stalled") + stalled = true + } + prev = progress + + select { + case <-tick.C: + case <-doneCh: + // Generate a final status update. + m.log.Printf("sourced %d records in %s", m.progress.Check(), time.Since(startTime)) + return + } + } + }() + } + + m.newNexter = func(c int) (IDAllocator, error) { + var nexter IDAllocator + return nexter, nil + } + + onFinishRun = func() {} + + return onFinishRun, nil +} + +func (m *Main) setup() (onFinishRun func(), err error) { if err := m.validate(); err != nil { return nil, errors.Wrap(err, "validating configuration") } @@ -2045,12 +2151,35 @@ func (m *Main) batchFromSchema(schema []Field) ([]Recordizer, pilosabatch.Record return recordizers, batch, row, lookupWriteIdxs, nil } -func (m *Main) newBatch(fields []*pilosaclient.Field) (pilosabatch.RecordBatch, error) { +func (m *Main) newBatch(clientFields []*pilosaclient.Field) (pilosabatch.RecordBatch, error) { + cfg := pilosabatch.Config{ + Size: m.BatchSize, + MaxStaleness: m.BatchMaxStaleness, + } + + // Table. + ii := pilosaclient.FromClientIndex(m.index) + tbl := pilosacore.IndexInfoToTable(ii) + + // Fields. + fields := pilosaclient.FromClientFields(clientFields) + + // If a custom Batcher has been defined, use that. Otherwise default to + // using the standard featurebase batch. + if m.Batcher != nil { + return m.Batcher.NewBatch(cfg, tbl, pilosacore.FieldInfosToFields(fields)) + } + + return m.newFeaturebaseBatch(cfg, tbl, fields) +} + +// newFeaturebaseBatch returns a featurebase.Batch based on the provided fields. +func (m *Main) newFeaturebaseBatch(cfg pilosabatch.Config, tbl *dax.Table, fields []*pilosacore.FieldInfo) (pilosabatch.RecordBatch, error) { opts := []pilosabatch.BatchOption{ pilosabatch.OptLogger(m.log), pilosabatch.OptCacheMaxAge(m.CacheLength), pilosabatch.OptSplitBatchMode(m.ExpSplitBatchMode), - pilosabatch.OptMaxStaleness(m.BatchMaxStaleness), + pilosabatch.OptMaxStaleness(cfg.MaxStaleness), pilosabatch.OptKeyTranslateBatchSize(m.KeyTranslateBatchSize), pilosabatch.OptUseShardTransactionalEndpoint(m.UseShardTransactionalEndpoint), } @@ -2063,10 +2192,7 @@ func (m *Main) newBatch(fields []*pilosaclient.Field) (pilosabatch.RecordBatch, } opts = append(opts, pilosabatch.OptImporter(importer)) - ii := pilosaclient.FromClientIndex(m.index) - tbl := pilosacore.IndexInfoToTable(ii) - - return pilosabatch.NewBatch(importer, m.BatchSize, tbl, pilosaclient.FromClientFields(fields), opts...) + return pilosabatch.NewBatch(importer, cfg.Size, tbl, fields, opts...) } // validateField ensures that the field is configured correctly. diff --git a/idk/interfaces.go b/idk/interfaces.go index a721694c6..07b4bcb1a 100644 --- a/idk/interfaces.go +++ b/idk/interfaces.go @@ -1374,7 +1374,7 @@ func (n *nopSchemaManager) FinishTransaction(id string) (*pilosacore.Transaction return nil, nil } func (n *nopSchemaManager) Schema() (*pilosaclient.Schema, error) { - return nil, nil + return pilosaclient.NewSchema(), nil } func (n *nopSchemaManager) SyncIndex(index *pilosaclient.Index) error { return nil diff --git a/idk/kafka/source_test.go b/idk/kafka/source_test.go index d117b83e9..de95c20a7 100644 --- a/idk/kafka/source_test.go +++ b/idk/kafka/source_test.go @@ -17,7 +17,6 @@ import ( "testing" "time" - "github.com/confluentinc/confluent-kafka-go/kafka" confluent "github.com/confluentinc/confluent-kafka-go/kafka" "github.com/featurebasedb/featurebase/v3/idk" "github.com/featurebasedb/featurebase/v3/idk/common" @@ -916,14 +915,14 @@ func tPutRecordsKafka(t *testing.T, p *confluent.Producer, topic string, schemaI func tPutRecordsKafkaPartition(t *testing.T, p *confluent.Producer, topic string, schemaID int, schema *liavro.Codec, key string, partition int32, records ...map[string]interface{}) { t.Helper() - delivery_chan := make(chan kafka.Event, 10000) + delivery_chan := make(chan confluent.Event, 10000) for _, record := range records { data, err := endcodeAvro(schemaID, schema, record) if err != nil { t.Fatalf("encoding record: %v", err) } - err = p.Produce(&kafka.Message{ - TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: partition}, + err = p.Produce(&confluent.Message{ + TopicPartition: confluent.TopicPartition{Topic: &topic, Partition: partition}, Key: []byte(key), Value: data, }, delivery_chan) @@ -931,7 +930,7 @@ func tPutRecordsKafkaPartition(t *testing.T, p *confluent.Producer, topic string t.Fatalf("producing record: %v", err) } e := <-delivery_chan - m := e.(*kafka.Message) + m := e.(*confluent.Message) if m.TopicPartition.Error != nil { t.Fatalf("Delivery failed: %v\n", m.TopicPartition.Error) diff --git a/idk/kafka_static/source.go b/idk/kafka_static/source.go index b8c871988..6f89f750c 100644 --- a/idk/kafka_static/source.go +++ b/idk/kafka_static/source.go @@ -26,14 +26,23 @@ import ( // source. It is not threadsafe! Due to the way Kafka clients work, to // achieve concurrency, create multiple Sources. type Source struct { - Hosts []string - Topics []string - Group string - TLS idk.TLSConfig - Log logger.Logger - Timeout time.Duration - SkipOld bool - Header string + Hosts []string + Topics []string + Group string + TLS idk.TLSConfig + Log logger.Logger + Timeout time.Duration + SkipOld bool + + // Header is a file or url referencing a file containing JSON header + // configuration. + Header string + + // HeaderFields can be provided instead of Header. It is a slice of + // RawFields which will be marshalled and parsed the same way a JSON object + // in Header would be. It is used only if a Header is not provided. + HeaderFields []idk.RawField + S3Region string AllowMissingFields bool @@ -162,24 +171,31 @@ func (r *Record) Data() []interface{} { // Open initializes the kafka source. func (s *Source) Open() error { - if len(s.Header) == 0 { - return errors.New("needs header specification file") + if len(s.Header) == 0 && len(s.HeaderFields) == 0 { + return errors.New("needs header specification (file or fields)") } - { - headerData, err := s.readFileOrURL(s.Header) + var headerData []byte + var err error + if s.Header != "" { + headerData, err = s.readFileOrURL(s.Header) if err != nil { return errors.Wrap(err, "reading header file") } - - schema, paths, err := idk.ParseHeader(headerData) + } else { + headerData, err = json.Marshal(s.HeaderFields) if err != nil { - return errors.Wrap(err, "processing header") + return errors.Wrap(err, "marshalling header fields") } - s.schema = schema - s.paths = paths } + schema, paths, err := idk.ParseHeader(headerData) + if err != nil { + return errors.Wrap(err, "processing header") + } + s.schema = schema + s.paths = paths + // init (custom) config, enable errors and notifications config := segmentio.ReaderConfig{ Brokers: s.Hosts, diff --git a/schema.go b/schema.go index e78a12f18..29efd4b29 100644 --- a/schema.go +++ b/schema.go @@ -280,6 +280,15 @@ func featurebaseFieldOptionsToEpoch(fo *FieldOptions) time.Time { return time.Unix(0, epochNano) } +// FieldInfosToFields converts a []*featurebase.FieldInfo to a []*dax.Field. +func FieldInfosToFields(fis []*FieldInfo) []*dax.Field { + fs := make([]*dax.Field, 0, len(fis)) + for i := range fis { + fs = append(fs, FieldInfoToField(fis[i])) + } + return fs +} + // // Functions to convert from dax to featurebase. //