mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
Merge branch 'jrb/cli-kafka-sasl' into jrb/cli-kafka-delete
This commit is contained in:
commit
04e029d36c
36 changed files with 1601 additions and 152 deletions
|
|
@ -192,13 +192,30 @@ build fbsql amd64:
|
|||
script:
|
||||
- export SOURCE_DATE_EPOCH=$(git log -1 --pretty=%ct)
|
||||
- date
|
||||
- GOOS="linux" GOARCH="amd64" make docker-build-fbsql BUILD_CGO=1
|
||||
- GOOS="linux" GOARCH="amd64" make docker-build-fbsql
|
||||
- GOOS="darwin" GOARCH="amd64" make docker-build-fbsql
|
||||
artifacts:
|
||||
paths:
|
||||
- ./build/fbsql_*
|
||||
|
||||
build fbsql arm64:
|
||||
build fbsql arm64 darwin:
|
||||
stage: test
|
||||
variables:
|
||||
BUILD_NAME: build_${CI_COMMIT_SHA}_${CI_CONCURRENT_ID}
|
||||
tags:
|
||||
- shell
|
||||
rules:
|
||||
- if: '$CI_PIPELINE_SOURCE == "push" || $CI_PIPELINE_SOURCE == "schedule" || $CI_PIPELINE_SOURCE == "web"'
|
||||
script:
|
||||
- export SOURCE_DATE_EPOCH=$(git log -1 --pretty=%ct)
|
||||
- date
|
||||
- GOOS="darwin" GOARCH="arm64" make docker-build-fbsql
|
||||
|
||||
artifacts:
|
||||
paths:
|
||||
- ./build/fbsql_*
|
||||
|
||||
build fbsql arm64 linux:
|
||||
stage: test
|
||||
variables:
|
||||
BUILD_NAME: build_${CI_COMMIT_SHA}_${CI_CONCURRENT_ID}
|
||||
|
|
@ -209,8 +226,8 @@ build fbsql arm64:
|
|||
script:
|
||||
- export SOURCE_DATE_EPOCH=$(git log -1 --pretty=%ct)
|
||||
- date
|
||||
- GOOS="linux" GOARCH="arm64" make docker-build-fbsql BUILD_CGO=1
|
||||
- GOOS="darwin" GOARCH="arm64" make docker-build-fbsql
|
||||
- GOOS="linux" GOARCH="arm64" make docker-build-fbsql
|
||||
|
||||
artifacts:
|
||||
paths:
|
||||
- ./build/fbsql_*
|
||||
|
|
@ -272,6 +289,7 @@ run go tests race:
|
|||
- echo "Running featurebase race tests..."
|
||||
- PKG_LIST=$(go list ./... | grep -Ev 'internal/clustertests|simulacraData|batch|idk|v3/dax/test/dax' | paste -s -d, -)
|
||||
- export TMPDIR=/mnt/ramdisk/test-$CI_JOB_ID
|
||||
- export SKIP_INTEGRATION_TEST=true
|
||||
- mkdir -p $TMPDIR
|
||||
- go test -race -v -timeout=10m ${PKG_LIST//,/ }
|
||||
after_script:
|
||||
|
|
@ -302,6 +320,7 @@ run go tests:
|
|||
- echo "Running featurebase unit tests..."
|
||||
- PKG_LIST=$(go list ./... | grep -Ev 'internal/clustertests|simulacraData|batch|idk|v3/dax/test/dax' | paste -s -d, -)
|
||||
- export TMPDIR=/mnt/ramdisk/test-$CI_JOB_ID
|
||||
- export SKIP_INTEGRATION_TEST=true
|
||||
- mkdir -p $TMPDIR
|
||||
- go test -tags=shardwidth22 -timeout=10m -coverprofile=coverage.out -covermode=atomic -coverpkg=${PKG_LIST} ${PKG_LIST//,/ }
|
||||
after_script:
|
||||
|
|
@ -342,6 +361,37 @@ run go tests dax/test/dax:
|
|||
paths:
|
||||
- coverage-dax-integration.out
|
||||
|
||||
# fbsql test
|
||||
run go tests cli kafka integration:
|
||||
stage: integration
|
||||
image: golang:$GOVERSION
|
||||
tags:
|
||||
- shell
|
||||
- aws
|
||||
variables:
|
||||
KAFKA_RUNNER_TEST_FEATUREBASE_HOST: pilosa:10101
|
||||
KAFKA_RUNNER_TEST_FEATUREBASEGRPC_HOST: pilosa:20101
|
||||
KAFKA_RUNNER_TEST_KAFKA_HOST: kafka:9092
|
||||
KAFKA_RUNNER_TEST_REGISTRY_HOST: schema-registry:8081
|
||||
rules:
|
||||
- if: '$CI_PIPELINE_SOURCE == "push" || $CI_PIPELINE_SOURCE == "schedule" || $CI_PIPELINE_SOURCE == "web"'
|
||||
script:
|
||||
- echo "running fbsql integration tests"
|
||||
- cd ./idk/
|
||||
- make test-cli
|
||||
after_script:
|
||||
- cd ./idk/
|
||||
- make save-pilosa-logs
|
||||
- make shutdown
|
||||
artifacts:
|
||||
paths:
|
||||
- ./idk/testdata/*_coverage.out
|
||||
needs:
|
||||
- job: build amd container fb
|
||||
- job: build fbsql amd64
|
||||
- job: build fbsql arm64 linux
|
||||
- job: build fbsql arm64 darwin
|
||||
|
||||
# idk tests
|
||||
run go tests idk race:
|
||||
variables:
|
||||
|
|
@ -471,16 +521,23 @@ upload to sonarcloud:
|
|||
package for linux amd64:
|
||||
stage: build
|
||||
image: golang:$GOVERSION
|
||||
extends: .go-cache
|
||||
rules:
|
||||
- if: '$CI_PIPELINE_SOURCE == "push" || $CI_PIPELINE_SOURCE == "schedule" || $CI_PIPELINE_SOURCE == "web"'
|
||||
variables:
|
||||
GOOS: "linux"
|
||||
GOARCH: "amd64"
|
||||
script:
|
||||
- export VERSION=$(git describe --tags 2>/dev/null || echo unknown)
|
||||
- echo 'deb [trusted=yes] https://repo.goreleaser.com/apt/ /' | tee /etc/apt/sources.list.d/goreleaser.list
|
||||
- apt update && apt install nfpm=2.11.3
|
||||
- make package
|
||||
- apt-get update -y && apt-get install -y -qq nfpm=2.11.3
|
||||
- nfpm version # smoke
|
||||
- mv "featurebase_${GOOS}_${GOARCH}" featurebase
|
||||
- mv "./build/fbsql_${GOOS}_${GOARCH}" fbsql
|
||||
- nfpm package --packager deb --target "featurebase.${VERSION}.${GOARCH}.deb"
|
||||
- nfpm package --packager rpm --target "featurebase.${VERSION}.${GOARCH}.rpm"
|
||||
needs:
|
||||
- build featurebase
|
||||
- build fbsql amd64
|
||||
artifacts:
|
||||
paths:
|
||||
- "*.deb"
|
||||
|
|
@ -509,16 +566,23 @@ trigger_m-cloud-images:
|
|||
package for linux arm64:
|
||||
stage: build
|
||||
image: golang:$GOVERSION
|
||||
extends: .go-cache
|
||||
rules:
|
||||
- if: '$CI_PIPELINE_SOURCE == "push" || $CI_PIPELINE_SOURCE == "schedule" || $CI_PIPELINE_SOURCE == "web"'
|
||||
variables:
|
||||
GOOS: "linux"
|
||||
GOARCH: "arm64"
|
||||
script:
|
||||
- export VERSION=$(git describe --tags 2>/dev/null || echo unknown)
|
||||
- echo 'deb [trusted=yes] https://repo.goreleaser.com/apt/ /' | tee /etc/apt/sources.list.d/goreleaser.list
|
||||
- apt update && apt install nfpm=2.11.3
|
||||
- make package
|
||||
- apt-get update -y && apt-get install -y -qq nfpm=2.11.3
|
||||
- nfpm version # smoke
|
||||
- mv "featurebase_${GOOS}_${GOARCH}" featurebase
|
||||
- mv "./build/fbsql_${GOOS}_${GOARCH}" fbsql
|
||||
- nfpm package --packager deb --target "featurebase.${VERSION}.${GOARCH}.deb"
|
||||
- nfpm package --packager rpm --target "featurebase.${VERSION}.${GOARCH}.rpm"
|
||||
needs:
|
||||
- build featurebase
|
||||
- build fbsql arm64 linux
|
||||
artifacts:
|
||||
paths:
|
||||
- "*.deb"
|
||||
|
|
@ -707,7 +771,8 @@ s3 dump:
|
|||
needs:
|
||||
- job: build featurebase
|
||||
- job: build fbsql amd64
|
||||
- job: build fbsql arm64
|
||||
- job: build fbsql arm64 linux
|
||||
- job: build fbsql arm64 darwin
|
||||
|
||||
s3 dump tag:
|
||||
stage: post build
|
||||
|
|
|
|||
78
Dockerfile-darwin-cgo-builder
Normal file
78
Dockerfile-darwin-cgo-builder
Normal file
|
|
@ -0,0 +1,78 @@
|
|||
# Dockerfile for building a container which can cross compile darwin with cgo
|
||||
# on linux
|
||||
#
|
||||
# fbsql depends https://github.com/confluentinc/confluent-kafka-go thus
|
||||
# requiring cgo to build. This generally isn't an issue unless you want to cross
|
||||
# compile the darwin build on linux. To cross compile darwin build on linux with
|
||||
# cgo, you need a C cross compiler. This Dockerfile utilizes
|
||||
# "https://github.com/tpoechtrager/osxcross.git for this.
|
||||
#
|
||||
# In order for osxcross to build compilers, it requires Xcode. Xcode is a
|
||||
# complete developer toolset for creating apps for Mac, iPhone, etc. The .xip
|
||||
# file required for this docker files can be founder here:
|
||||
# https://developer.apple.com/download/all/. Decide which Xcode version you'd
|
||||
# like to build with AND what version of go you'd like to build with.
|
||||
#
|
||||
|
||||
# Usage
|
||||
#
|
||||
# Assumptions
|
||||
# 1. For this dockerfile to work, you need to be serving the Xcode*.xip file.
|
||||
# Here is the command I ran from the directory that contained the
|
||||
# Xcode*.xip I had. Notice my private IP is hardcoded below. You'll need to
|
||||
# update that below to you private IP. This was done (as opposed to using
|
||||
# COPY) to keep the final image as small as possible.
|
||||
#
|
||||
# python3 -m http.server --bind 192.168.1.212
|
||||
#
|
||||
# 2. The version of Xcode and the underlying sdk are hard coded below. If
|
||||
# you're not using Xcode13.4.1.xip, you'll need to change both version
|
||||
# below.
|
||||
#
|
||||
# 3. Note the architecture used to build this build may make it easier /
|
||||
# harder to build what you intent to build.
|
||||
#
|
||||
# Here is a build example:
|
||||
#
|
||||
# docker build -f Dockerfile-darwin-cgo-builder -t <image_name>:<tag> .
|
||||
|
||||
|
||||
FROM ubuntu:focal AS builder
|
||||
|
||||
ARG DEBIAN_FRONTEND=noninteractive
|
||||
|
||||
WORKDIR /
|
||||
RUN apt-get update -y -qq && apt-get install -y -qq \
|
||||
build-essential \
|
||||
git \
|
||||
musl-tools \
|
||||
netcat \
|
||||
unixodbc \
|
||||
unixodbc-dev \
|
||||
clang \
|
||||
libxml2-dev \
|
||||
liblzma-dev \
|
||||
cmake \
|
||||
cpio \
|
||||
libssl-dev \
|
||||
zlib1g-dev \
|
||||
libbz2-dev \
|
||||
wget \
|
||||
libmpc-dev \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
WORKDIR /
|
||||
RUN git clone https://github.com/tpoechtrager/osxcross.git \
|
||||
&& wget http://192.168.1.212:8000/Xcode_13.4.1.xip \
|
||||
&& cd /osxcross \
|
||||
&& ./tools/gen_sdk_package_pbzx.sh /Xcode_13.4.1.xip \
|
||||
&& mv ./MacOSX12.3.sdk.tar.xz /osxcross/tarballs \
|
||||
&& rm /Xcode_13.4.1.xip
|
||||
|
||||
WORKDIR /osxcross
|
||||
RUN export UNATTENDED=1 \
|
||||
&& export PATH=/osxcross/build/cctools-port:$PATH \
|
||||
&& export PATH=/osxcross/build/cctools-port/cctools:$PATH \
|
||||
&& export PATH=/osxcross/build/cctools-port/cctools:$PATH \
|
||||
&& export PATH=/osxcross/build/:$PATH \
|
||||
&& ./build.sh \
|
||||
82
Dockerfile-fbsql-darwin
Normal file
82
Dockerfile-fbsql-darwin
Normal file
|
|
@ -0,0 +1,82 @@
|
|||
# Dockerfile for building fbsql on darwin
|
||||
#
|
||||
# fbsql depends https://github.com/confluentinc/confluent-kafka-go thus
|
||||
# requiring cgo to build. This generally isn't an issue unless you want to cross
|
||||
# compile the darwin build on linux. To cross compile darwin build on linux with
|
||||
# cgo, you need a C cross compiler. This Dockerfile utilizes
|
||||
# "https://github.com/tpoechtrager/osxcross.git for this.
|
||||
#
|
||||
# In order for osxcross to build compilers, it requires Xcode which is a
|
||||
# complete developer toolset for creating apps for Mac, iPhone, etc. This
|
||||
# version of the fbsql darwin builder does this build from an Xcode.xip file
|
||||
# which is more time consuming that doing it from an sdk tarball.
|
||||
#
|
||||
# For this to work, that Xcode.xip file must be in the working directory. See
|
||||
# the XCODE argument below. It's currently hardcoded but could be passes as a
|
||||
# parameter if needed (see the MAKE_FLAGS argument for example).
|
||||
|
||||
ARG GO_VERSION=latest
|
||||
|
||||
FROM jacobbrinlee/darwin-cgo-builder:13.4.1 AS builder
|
||||
|
||||
WORKDIR /
|
||||
RUN apt-get update -y -qq && apt-get install -y -qq \
|
||||
build-essential \
|
||||
git \
|
||||
musl-tools \
|
||||
netcat \
|
||||
unixodbc \
|
||||
unixodbc-dev \
|
||||
clang \
|
||||
libxml2-dev \
|
||||
liblzma-dev \
|
||||
cmake \
|
||||
cpio \
|
||||
libssl-dev \
|
||||
zlib1g-dev \
|
||||
libbz2-dev \
|
||||
pkg-config \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
RUN ["git", "clone", "https://github.com/edenhill/librdkafka.git"]
|
||||
WORKDIR /librdkafka
|
||||
RUN ./configure --prefix /usr &&\
|
||||
make && \
|
||||
make install
|
||||
|
||||
WORKDIR /featurebase
|
||||
|
||||
COPY . .
|
||||
|
||||
ARG MAKE_FLAGS
|
||||
ARG GO_BUILD_FLAGS
|
||||
ARG SOURCE_DATE_EPOCH
|
||||
|
||||
WORKDIR /featurebase/
|
||||
|
||||
COPY --from=golang:1.19 /usr/local/go /usr/local/go
|
||||
|
||||
ENV PATH="/usr/local/go/bin:${PATH}"
|
||||
ENV PATH="/osxcross/build/:${PATH}"
|
||||
ENV PATH="/osxcross/build/cctools-port/:${PATH}"
|
||||
ENV PATH="/osxcross/build/cctools-port/cctools:${PATH}"
|
||||
|
||||
ENV SOURCE_DATE_EPOCH=${SOURCE_DATE_EPOCH}
|
||||
RUN make build-fbsql-darwin GO_BUILD_FLAGS="-mod=vendor ${GO_BUILD_FLAGS}" ${MAKE_FLAGS}
|
||||
|
||||
FROM ubuntu:jammy AS runner
|
||||
|
||||
RUN apt-get update -y -qq && apt-get install -y -qq \
|
||||
ca-certificates \
|
||||
musl-tools \
|
||||
netcat \
|
||||
unixodbc-dev \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
COPY --from=builder /featurebase/fbsql /usr/local/bin/
|
||||
|
||||
# Verify that the linker can find everything.
|
||||
FROM runner AS linkcheck
|
||||
RUN if [ -e /usr/local/bin/fbsql ] ; then ldd /usr/local/bin/fbsql; fi
|
||||
|
||||
FROM runner
|
||||
|
|
@ -1,6 +1,6 @@
|
|||
ARG GO_VERSION=1.19
|
||||
|
||||
FROM golang:1.19-buster as builder
|
||||
FROM golang:1.19-buster AS builder
|
||||
|
||||
WORKDIR /
|
||||
RUN apt-get update -y -qq && apt-get install -y -qq \
|
||||
|
|
@ -28,9 +28,9 @@ ARG SOURCE_DATE_EPOCH
|
|||
WORKDIR /featurebase/
|
||||
|
||||
ENV SOURCE_DATE_EPOCH=${SOURCE_DATE_EPOCH}
|
||||
RUN make build-fbsql GO_BUILD_FLAGS="-mod=vendor ${GO_BUILD_FLAGS}" ${MAKE_FLAGS}
|
||||
RUN make build-fbsql-linux GO_BUILD_FLAGS="-mod=vendor ${GO_BUILD_FLAGS}" ${MAKE_FLAGS}
|
||||
|
||||
FROM ubuntu:20.04 as runner
|
||||
FROM ubuntu:20.04 AS runner
|
||||
|
||||
RUN apt-get update -y -qq && apt-get install -y -qq \
|
||||
ca-certificates \
|
||||
|
|
@ -42,7 +42,7 @@ RUN apt-get update -y -qq && apt-get install -y -qq \
|
|||
COPY --from=builder /featurebase/fbsql /usr/local/bin/
|
||||
|
||||
# Verify that the linker can find everything.
|
||||
FROM runner as linkcheck
|
||||
FROM runner AS linkcheck
|
||||
RUN if [ -e /usr/local/bin/fbsql ] ; then ldd /usr/local/bin/fbsql; fi
|
||||
|
||||
FROM runner
|
||||
48
Makefile
48
Makefile
|
|
@ -123,7 +123,7 @@ build:
|
|||
|
||||
package:
|
||||
GOOS=$(GOOS) GOARCH=$(GOARCH) $(MAKE) build
|
||||
GOOS=$(GOOS) GOARCH=$(GOARCH) $(MAKE) build-fbsql
|
||||
GOOS=$(GOOS) GOARCH=$(GOARCH) $(MAKE) docker-build-fbsql
|
||||
GOARCH=$(GOARCH) VERSION=$(VERSION) nfpm package --packager deb --target featurebase.$(VERSION).$(GOARCH).deb
|
||||
GOARCH=$(GOARCH) VERSION=$(VERSION) nfpm package --packager rpm --target featurebase.$(VERSION).$(GOARCH).rpm
|
||||
|
||||
|
|
@ -366,38 +366,44 @@ BUILD_NAME ?= fbsql-build
|
|||
LDFLAGS_STATIC="-linkmode external -extldflags \"-static\" -X 'github.com/featurebasedb/featurebase/v3/fbsql.Version=$(VERSION)' -X 'github.com/featurebasedb/featurebase/v3/fbsql.BuildTime=$(BUILD_TIME)' "
|
||||
|
||||
UNAME_P := $(shell uname -p)
|
||||
BUILD_CGO ?= 0
|
||||
|
||||
# Build fbsql
|
||||
build-fbsql:
|
||||
@echo GOOS=$(GOOS) GOARCH=$(GOARCH) uname -p=$(UNAME_P) build_cgo=$(BUILD_CGO)
|
||||
ifeq ($(BUILD_CGO), 0)
|
||||
make build-fbsql-non-cgo
|
||||
endif
|
||||
ifeq ($(BUILD_CGO), 1)
|
||||
make build-fbsql-cgo
|
||||
ifeq ($(GOOS), linux)
|
||||
$(MAKE) build-fbsql-linux
|
||||
else ifeq ($(GOOS), darwin)
|
||||
$(MAKE) build-fbsql-darwin
|
||||
endif
|
||||
|
||||
build-fbsql-non-cgo:
|
||||
CGO_ENABLED=0 $(GO) build -ldflags $(LDFLAGS) $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
|
||||
build-fbsql-cgo:
|
||||
build-fbsql-linux:
|
||||
ifeq ($(GOARCH), arm64)
|
||||
CGO_ENABLED=1 $(GO) build -tags dynamic $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
env GOOS=linux GOARCH=arm64 CGO_ENABLED=1 $(GO) build -tags dynamic $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
else ifeq ($(GOARCH), amd64)
|
||||
env GOOS=linux GOARCH=amd64 CC=/usr/bin/musl-gcc CGO_ENABLED=1 $(GO) build -tags "musl static" -ldflags $(LDFLAGS_STATIC) $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
endif
|
||||
ifeq ($(GOARCH), amd64)
|
||||
CC=/usr/bin/musl-gcc CGO_ENABLED=1 $(GO) build -tags "musl static" -ldflags $(LDFLAGS_STATIC) $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
|
||||
build-fbsql-darwin:
|
||||
ifeq ($(GOARCH), arm64)
|
||||
env GOOS=darwin GOARCH=arm64 CGO_ENABLED=1 CC=/osxcross/target/bin/arm64-apple-darwin21.4-clang $(GO) build $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
else ifeq ($(GOARCH), amd64)
|
||||
env GOOS=darwin GOARCH=amd64 CGO_ENABLED=1 CC=/osxcross/target/bin/x86_64-apple-darwin21.4-clang $(GO) build $(GO_BUILD_FLAGS) -o fbsql ./cmd/fbsql
|
||||
endif
|
||||
|
||||
docker-build-fbsql: vendor
|
||||
ifeq ($(GOOS), linux)
|
||||
@$(eval export FBSQL_DOCKERFILE=Dockerfile-fbsql-linux)
|
||||
else ifeq ($(GOOS), darwin)
|
||||
@$(eval export FBSQL_DOCKERFILE=Dockerfile-fbsql-darwin)
|
||||
endif
|
||||
DOCKER_BUILDKIT=0 docker build \
|
||||
--file Dockerfile-fbsql \
|
||||
--build-arg GO_VERSION=$(GO_VERSION) \
|
||||
--build-arg MAKE_FLAGS="GOOS=$(GOOS) GOARCH=$(GOARCH) BUILD_CGO=$(BUILD_CGO)" \
|
||||
--build-arg GO_BUILD_FLAGS=$(GO_BUILD_FLAGS) \
|
||||
--build-arg SOURCE_DATE_EPOCH=$(SOURCE_DATE_EPOCH) \
|
||||
--target builder \
|
||||
--tag fbsql:$(BUILD_NAME) .
|
||||
--file=$(FBSQL_DOCKERFILE) \
|
||||
--build-arg GO_VERSION=$(GO_VERSION) \
|
||||
--build-arg MAKE_FLAGS="GOOS=$(GOOS) GOARCH=$(GOARCH)" \
|
||||
--build-arg GO_BUILD_FLAGS=$(GO_BUILD_FLAGS) \
|
||||
--build-arg SOURCE_DATE_EPOCH=$(SOURCE_DATE_EPOCH) \
|
||||
--target builder \
|
||||
--tag fbsql:$(BUILD_NAME) .
|
||||
mkdir -p build
|
||||
docker create --name $(BUILD_NAME) fbsql:$(BUILD_NAME)
|
||||
docker cp $(BUILD_NAME):/featurebase/fbsql ./build/fbsql_$(GOOS)_$(GOARCH)
|
||||
|
|
|
|||
|
|
@ -175,10 +175,13 @@ func buildBulkInsert(tbl *dax.Table, fields []*dax.Field, ids []interface{}, row
|
|||
sb.WriteString(strings.Join(flds, ","))
|
||||
|
||||
// MAP
|
||||
|
||||
// map values in the MAP clause
|
||||
keyType := dax.BaseTypeID
|
||||
if tbl.StringKeys() {
|
||||
keyType = dax.BaseTypeString
|
||||
}
|
||||
|
||||
sb.WriteString(`) MAP ('$._id' `)
|
||||
sb.WriteString(keyType)
|
||||
sb.WriteString(`,`)
|
||||
|
|
@ -191,9 +194,16 @@ func buildBulkInsert(tbl *dax.Table, fields []*dax.Field, ids []interface{}, row
|
|||
// 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]
|
||||
if keyType == dax.BaseTypeID {
|
||||
m[string(dax.PrimaryKeyFieldName)] = ids[i]
|
||||
} else {
|
||||
// ids for key index can be string or []byte
|
||||
m[string(dax.PrimaryKeyFieldName)] = fmt.Sprintf("%s", ids[i])
|
||||
}
|
||||
|
||||
// Write the rest of the data values.
|
||||
for col := range rows[i] {
|
||||
m[fmt.Sprintf("col_%d", col)] = rows[i][col]
|
||||
|
|
|
|||
516
cli/cli_kafka_integration_test.go
Normal file
516
cli/cli_kafka_integration_test.go
Normal file
|
|
@ -0,0 +1,516 @@
|
|||
package cli_test
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/cli"
|
||||
"github.com/featurebasedb/featurebase/v3/cli/kafka"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/idk/kafka/csrc"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
confluent "github.com/confluentinc/confluent-kafka-go/kafka"
|
||||
|
||||
avro "github.com/linkedin/goavro/v2"
|
||||
)
|
||||
|
||||
// Struct that contains the services required to run the kafka runner tests
|
||||
type KafkaRunnerTestServices struct {
|
||||
featurebaseHost string
|
||||
featurebaseGRPCHost string
|
||||
kafkaHost string
|
||||
registryHost string
|
||||
}
|
||||
|
||||
func envOr(envName, defaultVal string) string {
|
||||
if val, ok := os.LookupEnv(envName); !ok {
|
||||
return defaultVal
|
||||
} else {
|
||||
return val
|
||||
}
|
||||
}
|
||||
|
||||
func getKafkaRunnerTestServices() *KafkaRunnerTestServices {
|
||||
return &KafkaRunnerTestServices{
|
||||
featurebaseHost: envOr("KAFKA_RUNNER_TEST_FEATUREBASE_HOST", "localhost:10101"),
|
||||
featurebaseGRPCHost: envOr("KAFKA_RUNNER_TEST_FEATUREBASEGRPC_HOST", "localhost:20101"),
|
||||
kafkaHost: envOr("KAFKA_RUNNER_TEST_KAFKA_HOST", "localhost:9092"),
|
||||
registryHost: envOr("KAFKA_RUNNER_TEST_REGISTRY_HOST", "localhost:8081"),
|
||||
}
|
||||
}
|
||||
|
||||
// Struct used for TestRunner which contains the information needed to ingest
|
||||
// data to kafka, configure the runner, and test that the runner successfully
|
||||
// ran.
|
||||
type kafkaRunnerTest struct {
|
||||
ConfigFile string // path to the configuration file used for the runner
|
||||
DataFile string // path to the data file to populate kafka with
|
||||
CreateTableStmt string // statement used to create table prior to ingest
|
||||
SchemaFile string // path to schema file for schema encoded messages (e.g. avro)
|
||||
Tests []testQuery // list of test which are 2-tuples of query and expected results
|
||||
|
||||
}
|
||||
|
||||
// Struct which contains a query to run agaisnt featurebase and the expected
|
||||
// results from that query.
|
||||
type testQuery struct {
|
||||
Query string
|
||||
ExpectedResp string
|
||||
}
|
||||
|
||||
// A slice of KafkaRunnerTest structs that will be used in TestKafkaRunner test
|
||||
// function.
|
||||
var kafkaRunnerTests = []kafkaRunnerTest{
|
||||
{ // id keys json
|
||||
ConfigFile: "config00.toml",
|
||||
DataFile: "data00.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by name",
|
||||
ExpectedResp: `[[1,"a",20,["hob1","hob2"]],[2,"b",21,["hob2","hob3"]],[3,"c",22,["hob3","hob4"]],[4,"d",23,["hob4","hob5"]],[5,"e",24,["hob5","hob6"]],[6,"f",26,["hob6","hob7"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id ID, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
{ // string keys json
|
||||
ConfigFile: "config01.toml",
|
||||
DataFile: "data00.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by name",
|
||||
ExpectedResp: `[["1","a",20,["hob1","hob2"]],["2","b",21,["hob2","hob3"]],["3","c",22,["hob3","hob4"]],["4","d",23,["hob4","hob5"]],["5","e",24,["hob5","hob6"]],["6","f",26,["hob6","hob7"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id String, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
{ // two string keys json
|
||||
ConfigFile: "config02.toml",
|
||||
DataFile: "data00.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by name",
|
||||
ExpectedResp: `[["1|a","1","a",20,["hob1","hob2"]],["2|b","2","b",21,["hob2","hob3"]],["3|c","3","c",22,["hob3","hob4"]],["4|d","4","d",23,["hob4","hob5"]],["5|e","5","e",24,["hob5","hob6"]],["6|f","6","f",26,["hob6","hob7"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id String, id String, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
{ // string, id, and int json
|
||||
ConfigFile: "config03.toml",
|
||||
DataFile: "data00.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by name",
|
||||
ExpectedResp: `[["1|a|20",1,"a",20,["hob1","hob2"]],["2|b|21",2,"b",21,["hob2","hob3"]],["3|c|22",3,"c",22,["hob3","hob4"]],["4|d|23",4,"d",23,["hob4","hob5"]],["5|e|24",5,"e",24,["hob5","hob6"]],["6|f|26",6,"f",26,["hob6","hob7"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id String, id id, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
{ // missing values json
|
||||
ConfigFile: "config05.toml",
|
||||
DataFile: "data02.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by name",
|
||||
ExpectedResp: `[["3",null,22,["hob3","hob4"]],["1","a",20,["hob1","hob2"]],["2","b",21,["hob2","hob3"]],["4","d",null,["hob4","hob5"]],["5","e",24,null],["6","f",26,["hob6","hob7"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id String, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
{ // string keys avro
|
||||
ConfigFile: "config04.toml",
|
||||
DataFile: "data01.json",
|
||||
SchemaFile: "schema01.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by string_string",
|
||||
ExpectedResp: `[["h1iqc","58KIR","x5z8P",["iYeOV"],["eNKWF"],[255],4.41,"2023-02-19T08:52:56Z",[63,110,320,344,606],1676796776,null,["ASSAw"],["6TKzc","RKE3c","ZgkOB","eofzb","pjxqm"],[821],"2023-02-19T14:52:56Z",110,647,389,[29,257,289,388,606],0.40,["bool_bool"],["5HIn2","7EYSp","BmvHF","Qylqq","yTeUQ"],148,2.84,["5ptDx"]],["yg8hY","5HIn2","qK5TE",["byHh9"],["u2Yr4"],[839],3.23,"2023-01-30T06:56:05Z",[63,582,629,680,690],1675061765,null,["911oj"],["d0U7s","dxKKn","fjQK2","m5d59","nVQrd"],[533],"2023-01-30T12:56:05Z",433,809,168,[115,172,175,257,969],1.27,["bool_bool"],["F0uC4","KMZnH","OKNV2","VBcyJ","wNZ7o"],680,1156.06,["tvNOB"]],["DY2Ui","8MGwy","vTwn4",["pjxqm"],["DDLN5"],[984],4.23,"2023-02-12T18:37:16Z",[113,733,751,772,975],1676227036,null,["tyP3m"],["8MGwy","XzEHj","gjWEI","v31XN","xE5jX"],[931],"2023-02-13T00:37:16Z",63,430,297,[72,297,384,694,898],0.83,["bool_bool"],["d0U7s","sDdtS","u2Yr4","y2Y7b"],388,1.26,["kUbdU"]],["tElMR","FW39I","FW39I",["n9HUP"],["PNB4s"],[289],2.19,"2023-02-20T17:04:21Z",[2,289,389,680,958],1676912661,null,["58KIR"],["58KIR","6TKzc","8MGwy","X9jWC"],[791],"2023-02-20T23:04:21Z",289,695,821,[102,220,387,606,890],2.65,["bool_bool"],["BmvHF","PNB4s","TLaUE","eofzb","vhisL"],2,0.95,["ARlcJ"]],["BmvHF","I1gXJ","thuky",["6TKzc"],["gjWEI"],[166],2.91,"2023-01-31T11:11:30Z",[284,289,388,890,975],1675163490,["bool_bool"],["X9jWC"],["5ptDx","Chgzr","EyQoi","TLaUE","tyP3m"],[232],"2023-01-31T17:11:30Z",857,320,286,[322,614,865,884,931],2.84,["bool_bool"],["F0uC4","VQs7y","byHh9","d0U7s","h1iqc"],879,500.86,["798ka"]],["ASSAw","LBTEU","EyQoi",["oxjI0"],["5ptDx"],[484],4.97,"2023-02-22T14:32:23Z",[168,399,639,792,809],1677076343,["bool_bool"],["iYeOV"],["XzEHj","iYeOV","rrkYB","uirDR","v31XN"],[322],"2023-02-22T20:32:23Z",533,23,320,[23,293,358,606,821],4.32,["bool_bool"],["PYE8V","X9jWC","vTwn4","x5z8P"],884,2.97,["kauLy"]],["RKE3c","TLaUE","YdwQY",["RKE3c"],["dxKKn"],[39],2.72,"2023-02-16T20:09:13Z",[63,172,220,358,857],1676578153,["bool_bool"],["dF6kx"],["5HIn2","KdTtE","nVQrd","wNZ7o","x5z8P"],[113],"2023-02-17T02:09:13Z",582,665,681,[220,647,665,731,778],1.36,["bool_bool"],["5HIn2","I6NST","Qylqq","gjWEI","tyP3m"],690,1156.01,["6TKzc"]],["u2Yr4","ZgkOB","6iGIm",["x5z8P"],["qK5TE"],[148],2.29,"2023-02-03T16:19:37Z",[63,148,839,958,984],1675441177,null,["7EYSp"],["Chgzr","DY2Ui","PYE8V","VBcyJ","u2Yr4"],[890],"2023-02-03T22:19:37Z",13,115,39,[13,167,629,731,772],2.93,["bool_bool"],["KdTtE","MVNow","YdwQY","aQQxr","kUbdU"],969,498.18,["sHaUv"]],["6TKzc","n9HUP","5HIn2",["h1iqc"],["t5f7R"],[72],0.78,"2023-02-23T05:04:34Z",[322,399,730,969,975],1677128674,["bool_bool"],["eofzb"],["BmvHF","C6xxn","PYE8V","xE5jX","yg8hY"],[676],"2023-02-23T11:04:34Z",430,387,797,[242,289,778,797,958],3.35,["bool_bool"],["6TKzc","jVVfZ","pjxqm","vK0WD","xE5jX"],23,1156.06,["YKLk9"]],["9z4aw","uirDR","BmvHF",["CKs1F"],["gL2Hg"],[647],0.95,"2023-02-16T07:53:59Z",[167,230,344,442,733],1676534039,null,["7EYSp"],["9z4aw","VQs7y","aQQxr","h1iqc","vbbuf"],[898],"2023-02-16T13:53:59Z",584,792,63,[284,344,394,442,614],3.23,["bool_bool"],["RPGAm","ZgkOB","iYeOV","tvNOB","u2Yr4"],344,1155.95,["5ptDx"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id string, string_string string, string_bytes string, pk2 stringset, stringset_bytes stringset, idset_long idset, decimal_double decimal(2), timestamp_bytes_ts timestamp, idset_longarray idset, dateint_bytes_ts int, bools stringset, stringset_string stringset, stringset_stringarray stringset, idset_int idset, timestamp_bytes_int timestamp, int_long int, id_long id, id_int id, idset_intarray idset, decimal_float decimal(2), bools-exists stringset, stringset_bytesarray stringset, int_int int, decimal_bytes decimal(2), pk1 stringset)",
|
||||
},
|
||||
{ // id keys avro
|
||||
ConfigFile: "config06.toml",
|
||||
DataFile: "data03.json",
|
||||
SchemaFile: "schema02.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by string_string",
|
||||
ExpectedResp: `[[14,"58KIR",[110,320,606],["6TKzc","RKE3c","ZgkOB","eofzb","pjxqm"],148,null,["bool_bool"]],[10,"5HIn2",[582,629,680],["d0U7s","dxKKn","fjQK2","m5d59","nVQrd"],680,null,["bool_bool"]],[16,"8MGwy",[113,733,975],["8MGwy","XzEHj","gjWEI","v31XN","xE5jX"],388,null,["bool_bool"]],[6,"FW39I",[201,680,958],["58KIR","6TKzc","8MGwy","X9jWC"],212,null,["bool_bool"]],[4,"I1gXJ",[284,890,975],["5ptDx","Chgzr","EyQoi","TLaUE","tyP3m"],879,["bool_bool"],["bool_bool"]],[2,"LBTEU",[168,792,809],["XzEHj","iYeOV","rrkYB","uirDR","v31XN"],884,["bool_bool"],["bool_bool"]],[8,"TLaUE",[172,630,857],["5HIn2","KdTtE","nVQrd","wNZ7o","x5z8P"],690,["bool_bool"],["bool_bool"]],[18,"ZgkOB",[148,635,839],["Chgzr","DY2Ui","PYE8V","VBcyJ","u2Yr4"],969,null,["bool_bool"]],[12,"n9HUP",[322,399,975],["BmvHF","C6xxn","PYE8V","xE5jX","yg8hY"],230,["bool_bool"],["bool_bool"]],[0,"uirDR",[167,230,442],["9z4aw","VQs7y","aQQxr","h1iqc","vbbuf"],344,null,["bool_bool"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id id, string_string string, idset_longarray idset, stringset_stringarray stringset, int_int int, bools stringset, bools-exists stringset)",
|
||||
},
|
||||
{ // compound keys avro
|
||||
ConfigFile: "config07.toml",
|
||||
DataFile: "data01.json",
|
||||
SchemaFile: "schema01.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "select * from <TABLE> order by string_string",
|
||||
ExpectedResp: `[["h1iqc|5ptDx|148","58KIR","x5z8P",["iYeOV"],["eNKWF"],[255],4.41,"2023-02-19T08:52:56Z",[63,110,320,344,606],1676796776,null,["ASSAw"],["6TKzc","RKE3c","ZgkOB","eofzb","pjxqm"],[821],"2023-02-19T14:52:56Z",110,647,389,[29,257,289,388,606],0.40,["bool_bool"],["5HIn2","7EYSp","BmvHF","Qylqq","yTeUQ"],148,2.84,["5ptDx"],["h1iqc"]],["yg8hY|tvNOB|680","5HIn2","qK5TE",["byHh9"],["u2Yr4"],[839],3.23,"2023-01-30T06:56:05Z",[63,582,629,680,690],1675061765,null,["911oj"],["d0U7s","dxKKn","fjQK2","m5d59","nVQrd"],[533],"2023-01-30T12:56:05Z",433,809,168,[115,172,175,257,969],1.27,["bool_bool"],["F0uC4","KMZnH","OKNV2","VBcyJ","wNZ7o"],680,1156.06,["tvNOB"],["yg8hY"]],["DY2Ui|kUbdU|388","8MGwy","vTwn4",["pjxqm"],["DDLN5"],[984],4.23,"2023-02-12T18:37:16Z",[113,733,751,772,975],1676227036,null,["tyP3m"],["8MGwy","XzEHj","gjWEI","v31XN","xE5jX"],[931],"2023-02-13T00:37:16Z",63,430,297,[72,297,384,694,898],0.83,["bool_bool"],["d0U7s","sDdtS","u2Yr4","y2Y7b"],388,1.26,["kUbdU"],["DY2Ui"]],["tElMR|ARlcJ|2","FW39I","FW39I",["n9HUP"],["PNB4s"],[289],2.19,"2023-02-20T17:04:21Z",[2,289,389,680,958],1676912661,null,["58KIR"],["58KIR","6TKzc","8MGwy","X9jWC"],[791],"2023-02-20T23:04:21Z",289,695,821,[102,220,387,606,890],2.65,["bool_bool"],["BmvHF","PNB4s","TLaUE","eofzb","vhisL"],2,0.95,["ARlcJ"],["tElMR"]],["BmvHF|798ka|879","I1gXJ","thuky",["6TKzc"],["gjWEI"],[166],2.91,"2023-01-31T11:11:30Z",[284,289,388,890,975],1675163490,["bool_bool"],["X9jWC"],["5ptDx","Chgzr","EyQoi","TLaUE","tyP3m"],[232],"2023-01-31T17:11:30Z",857,320,286,[322,614,865,884,931],2.84,["bool_bool"],["F0uC4","VQs7y","byHh9","d0U7s","h1iqc"],879,500.86,["798ka"],["BmvHF"]],["ASSAw|kauLy|884","LBTEU","EyQoi",["oxjI0"],["5ptDx"],[484],4.97,"2023-02-22T14:32:23Z",[168,399,639,792,809],1677076343,["bool_bool"],["iYeOV"],["XzEHj","iYeOV","rrkYB","uirDR","v31XN"],[322],"2023-02-22T20:32:23Z",533,23,320,[23,293,358,606,821],4.32,["bool_bool"],["PYE8V","X9jWC","vTwn4","x5z8P"],884,2.97,["kauLy"],["ASSAw"]],["RKE3c|6TKzc|690","TLaUE","YdwQY",["RKE3c"],["dxKKn"],[39],2.72,"2023-02-16T20:09:13Z",[63,172,220,358,857],1676578153,["bool_bool"],["dF6kx"],["5HIn2","KdTtE","nVQrd","wNZ7o","x5z8P"],[113],"2023-02-17T02:09:13Z",582,665,681,[220,647,665,731,778],1.36,["bool_bool"],["5HIn2","I6NST","Qylqq","gjWEI","tyP3m"],690,1156.01,["6TKzc"],["RKE3c"]],["u2Yr4|sHaUv|969","ZgkOB","6iGIm",["x5z8P"],["qK5TE"],[148],2.29,"2023-02-03T16:19:37Z",[63,148,839,958,984],1675441177,null,["7EYSp"],["Chgzr","DY2Ui","PYE8V","VBcyJ","u2Yr4"],[890],"2023-02-03T22:19:37Z",13,115,39,[13,167,629,731,772],2.93,["bool_bool"],["KdTtE","MVNow","YdwQY","aQQxr","kUbdU"],969,498.18,["sHaUv"],["u2Yr4"]],["6TKzc|YKLk9|23","n9HUP","5HIn2",["h1iqc"],["t5f7R"],[72],0.78,"2023-02-23T05:04:34Z",[322,399,730,969,975],1677128674,["bool_bool"],["eofzb"],["BmvHF","C6xxn","PYE8V","xE5jX","yg8hY"],[676],"2023-02-23T11:04:34Z",430,387,797,[242,289,778,797,958],3.35,["bool_bool"],["6TKzc","jVVfZ","pjxqm","vK0WD","xE5jX"],23,1156.06,["YKLk9"],["6TKzc"]],["9z4aw|5ptDx|344","uirDR","BmvHF",["CKs1F"],["gL2Hg"],[647],0.95,"2023-02-16T07:53:59Z",[167,230,344,442,733],1676534039,null,["7EYSp"],["9z4aw","VQs7y","aQQxr","h1iqc","vbbuf"],[898],"2023-02-16T13:53:59Z",584,792,63,[284,344,394,442,614],3.23,["bool_bool"],["RPGAm","ZgkOB","iYeOV","tvNOB","u2Yr4"],344,1155.95,["5ptDx"],["9z4aw"]]]`,
|
||||
},
|
||||
},
|
||||
CreateTableStmt: "(_id string, string_string string, string_bytes string, pk2 stringset, stringset_bytes stringset, idset_long idset, decimal_double decimal(2), timestamp_bytes_ts timestamp, idset_longarray idset, dateint_bytes_ts int, bools stringset, stringset_string stringset, stringset_stringarray stringset, idset_int idset, timestamp_bytes_int timestamp, int_long int, id_long id, id_int id, idset_intarray idset, decimal_float decimal(2), bools-exists stringset, stringset_bytesarray stringset, int_int int, decimal_bytes decimal(2), pk1 stringset)",
|
||||
},
|
||||
}
|
||||
|
||||
// TestKafkaRunner takes as input a slice of kafkaRunnerTest structs.
|
||||
// For each kafkaRunnerTest, TestKafkaRunner:
|
||||
// 1. Creates and configues a new cli.Command
|
||||
// 2. Reads kafka messages from a data file
|
||||
// 3. Encodes that data
|
||||
// 4. Writes the data to kafka
|
||||
// 5. Runs the cli.Command
|
||||
// 6. Confirms that the data was written to FeatureBase as expected
|
||||
func TestKafkaRunner(t *testing.T) {
|
||||
if testing.Short() || os.Getenv("SKIP_INTEGRATION_TEST") == "true" {
|
||||
t.Skip("skipping integration test")
|
||||
}
|
||||
|
||||
// Get host and port pair for services required for test (e.g. kafka and
|
||||
// featurebase)
|
||||
services := getKafkaRunnerTestServices()
|
||||
|
||||
for _, test := range kafkaRunnerTests {
|
||||
|
||||
// define path to config and data files
|
||||
kafkaConfig := "./kafka/testdata/runner/config/" + test.ConfigFile
|
||||
dataFile := "./kafka/testdata/runner/data/" + test.DataFile
|
||||
schemaFile := "./kafka/testdata/runner/schema/" + test.SchemaFile
|
||||
|
||||
// read data that will go to kafka
|
||||
var records []map[string]interface{}
|
||||
records, err := recordsFromFile(dataFile)
|
||||
if err != nil {
|
||||
t.Errorf("reading raw messages from data file: %s", err)
|
||||
}
|
||||
|
||||
// copy config file and replace values as defined in findAndReplace
|
||||
var findAndReplace = map[string]string{
|
||||
"KAFKA_SERVICE": services.kafkaHost,
|
||||
"SCHEMA_REGISTRY_SERVICE": services.registryHost,
|
||||
"MAX_MESSAGES": strconv.Itoa(len(records)),
|
||||
}
|
||||
if err := createTempFindAndReplace(kafkaConfig, findAndReplace); err != nil {
|
||||
t.Fatalf("creating temp config file: %s", err)
|
||||
}
|
||||
defer os.Remove(kafkaConfig + ".tmp")
|
||||
|
||||
// create new command
|
||||
fbsql := cli.NewCommand(logger.StderrLogger)
|
||||
|
||||
fbsql.Config.Host = strings.Split(services.featurebaseHost, ":")[0]
|
||||
fbsql.Config.Port = strings.Split(services.featurebaseHost, ":")[1]
|
||||
fbsql.Config.KafkaConfig = "/bad/path" // need a path to be i
|
||||
fbsql.Run(context.Background()) // creates fbsql's Queryer which we can then use below
|
||||
fbsql.Config.KafkaConfig = kafkaConfig + ".tmp" // when fbsql comes back, add kafka config
|
||||
|
||||
config, err := kafka.ConfigFromFile(kafkaConfig + ".tmp")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// create table to write to (kafka runner does not currently create tables)
|
||||
createTable := strings.NewReader(fmt.Sprintf("CREATE TABLE %s %s;", config.Table, test.CreateTableStmt))
|
||||
wqr, err := fbsql.Queryer.Query("", "", createTable)
|
||||
if err != nil {
|
||||
t.Fatalf("creating table: %s", err)
|
||||
}
|
||||
|
||||
if wqr.Error != "" {
|
||||
t.Error(wqr.Error)
|
||||
}
|
||||
|
||||
// delete all tables after test (note this is after all tests not after each test)
|
||||
defer func() {
|
||||
dropTable := strings.NewReader(fmt.Sprintf("DROP TABLE %s;", config.Table))
|
||||
_, err = fbsql.Queryer.Query("", "", dropTable)
|
||||
if err != nil {
|
||||
t.Errorf("dropping table: %s", err)
|
||||
}
|
||||
}()
|
||||
|
||||
// write data to kafka
|
||||
var messages [][]byte
|
||||
switch encode := config.Encode; encode {
|
||||
case "json":
|
||||
messages, err = encodeJSONMessages(records)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
case "avro":
|
||||
// get avro schema as string
|
||||
schemaBytes, err := os.ReadFile(schemaFile)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
schema := string(schemaBytes)
|
||||
|
||||
// post the schema to schema registry
|
||||
schemaID, err := postSchema(schema, "kafka-runner-subject-"+time.Now().String(), services.registryHost)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// encode avro messages
|
||||
messages, err = encodeAvroMessages(records, schema, schemaID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
default:
|
||||
t.Errorf("unsupported encoding type")
|
||||
}
|
||||
|
||||
if err = createTopic(config.Topics[0], services.kafkaHost); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if err = produceMessages(config.Topics[0], services.kafkaHost, messages); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() {
|
||||
if err = deleteTopic(config.Topics[0], services.kafkaHost); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}()
|
||||
|
||||
if err = fbsql.Run(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// and check the results...
|
||||
for _, q := range test.Tests {
|
||||
query := strings.Replace(q.Query, "<TABLE>", config.Table, 1)
|
||||
if resp, err := fbsql.Queryer.Query("", "", strings.NewReader(query)); err != nil {
|
||||
t.Fatal(err)
|
||||
} else {
|
||||
verifyQueryReponse(t, resp, q.ExpectedResp)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// createTempFindAndReplace reads a file (source), replaces substrings specified
|
||||
// in mapping, and writes the output to the same path as souce, appending
|
||||
// ".tmp".
|
||||
func createTempFindAndReplace(source string, mapping map[string]string) error {
|
||||
//Read all the contents of the original file
|
||||
bytesRead, err := ioutil.ReadFile(source)
|
||||
if err != nil {
|
||||
return errors.Errorf("%s", err)
|
||||
}
|
||||
|
||||
stringRead := string(bytesRead)
|
||||
|
||||
for find, replace := range mapping {
|
||||
stringRead = strings.Replace(stringRead, find, replace, 1)
|
||||
}
|
||||
|
||||
//Copy all the contents to the desitination file
|
||||
if err = ioutil.WriteFile(source+".tmp", []byte(stringRead), 0755); err != nil {
|
||||
return errors.Errorf("%s", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
// verifyQueryResponse compares the JSON binary representation of the Data field
|
||||
// in an WireQueryResponse against some expected string. All lists in the Data
|
||||
// field of the WireQueryResponse are sorted and then serialized. The test that
|
||||
// calls this function fails if the query responses are different. In this case,
|
||||
// the diff is displayed.
|
||||
func verifyQueryReponse(t *testing.T, wqr *featurebase.WireQueryResponse, expectedQuery string) {
|
||||
// we need to sort slices so json compare is accurate
|
||||
var data = make([][]interface{}, len(wqr.Data))
|
||||
for i, line := range wqr.Data {
|
||||
var newline = make([]interface{}, len(line))
|
||||
for j, element := range line {
|
||||
switch newElement := element.(type) {
|
||||
case featurebase.StringSet:
|
||||
sort.Strings(newElement)
|
||||
newline[j] = newElement
|
||||
case featurebase.IDSet:
|
||||
sort.Slice(newElement, func(i, j int) bool { return newElement[i] < newElement[j] })
|
||||
newline[j] = newElement
|
||||
default:
|
||||
newline[j] = element
|
||||
}
|
||||
}
|
||||
data[i] = newline
|
||||
}
|
||||
js, err := json.Marshal(data)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// t.Fatal(string(js))
|
||||
require.JSONEq(t, expectedQuery, string(js))
|
||||
}
|
||||
|
||||
// recordsFromFile reads lines from a file and builds an in memory structure.
|
||||
// The file pointed to by pathToRecords must be formated as new line delimited
|
||||
// JSON.
|
||||
func recordsFromFile(pathToRecords string) (records []map[string]interface{}, err error) {
|
||||
var data map[string]interface{}
|
||||
|
||||
recordsFile, err := os.Open(pathToRecords)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("opening records file: %s", err)
|
||||
}
|
||||
defer recordsFile.Close()
|
||||
|
||||
s := bufio.NewScanner(recordsFile)
|
||||
for s.Scan() {
|
||||
err := json.Unmarshal(s.Bytes(), &data)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("unmarshal json: %s", err)
|
||||
}
|
||||
records = append(records, data)
|
||||
data = make(map[string]interface{})
|
||||
}
|
||||
return records, nil
|
||||
}
|
||||
|
||||
// encodeJSONMessages takes a slice of JSON objects and convert them to a slice
|
||||
// of byte slices which are the binary representation of JSON.
|
||||
func encodeJSONMessages(records []map[string]interface{}) ([][]byte, error) {
|
||||
messages := make([][]byte, len(records))
|
||||
for i, record := range records {
|
||||
encoded, err := json.Marshal(record)
|
||||
if err != nil {
|
||||
return nil, errors.Errorf("marshaling records to JSON: %s", err)
|
||||
}
|
||||
messages[i] = encoded
|
||||
}
|
||||
return messages, nil
|
||||
}
|
||||
|
||||
// encodeAvroMessages takes a slice of JSON objects and convert them to a slice
|
||||
// of byte slices which are the binary avro encoding based on schema with a
|
||||
// specific schemaID.
|
||||
func encodeAvroMessages(records []map[string]interface{}, schema string, schemaID int) ([][]byte, error) {
|
||||
|
||||
// get a thing which can encode a byte slice based on a schema
|
||||
avroEncoder, err := avro.NewCodec(schema)
|
||||
if err != nil {
|
||||
return nil, errors.Errorf("getting avro encoder: %s", err)
|
||||
}
|
||||
|
||||
messages := make([][]byte, len(records))
|
||||
for i, record := range records {
|
||||
buf := make([]byte, 5, 1000)
|
||||
buf[0] = 0
|
||||
binary.BigEndian.PutUint32(buf[1:], uint32(schemaID))
|
||||
buf, err = avroEncoder.BinaryFromNative(buf, record)
|
||||
if err != nil {
|
||||
return nil, errors.Errorf("avro encoding record: %s", err)
|
||||
}
|
||||
messages[i] = buf
|
||||
}
|
||||
return messages, nil
|
||||
}
|
||||
|
||||
// postSchema takes an avro schema and subject and posts it to schema registry
|
||||
// at a specific url.
|
||||
func postSchema(schema, subj, schemaRegistryURL string) (schemaID int, err error) {
|
||||
schemaClient := csrc.NewClient("http://"+schemaRegistryURL, nil, nil)
|
||||
resp, err := schemaClient.PostSubjects(subj, schema)
|
||||
if err != nil {
|
||||
return 0, errors.Wrap(err, "posting schema")
|
||||
}
|
||||
return resp.ID, nil
|
||||
}
|
||||
|
||||
// createTopic creates a kafka topic on the kafka server pointed to by
|
||||
// bootstrapServers
|
||||
func createTopic(topic string, bootstrapServers string) error {
|
||||
|
||||
// create admin client needed to create topic, defer closing
|
||||
ac, err := confluent.NewAdminClient(&confluent.ConfigMap{"bootstrap.servers": bootstrapServers})
|
||||
if err != nil {
|
||||
return errors.Errorf("creating admin client: %s", err)
|
||||
}
|
||||
defer ac.Close()
|
||||
|
||||
// create a topic
|
||||
var ts = []confluent.TopicSpecification{
|
||||
{
|
||||
Topic: topic,
|
||||
NumPartitions: 1,
|
||||
ReplicationFactor: 1,
|
||||
},
|
||||
}
|
||||
_, err = ac.CreateTopics(context.Background(), ts, nil)
|
||||
if err != nil {
|
||||
return errors.Errorf("creating topics: %s", err)
|
||||
}
|
||||
// caller must delete if needed
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// produceMessages produces messages to a kafka topic on the kafka server
|
||||
// pointed to by bootstrapServer. Messages should be byte slices.
|
||||
func produceMessages(topic string, bootstrapServers string, messages [][]byte) error {
|
||||
|
||||
// create producer, defer closing
|
||||
p, err := confluent.NewProducer(&confluent.ConfigMap{"bootstrap.servers": bootstrapServers})
|
||||
if err != nil {
|
||||
return errors.Errorf("creating producer: %s", err)
|
||||
}
|
||||
defer p.Close()
|
||||
|
||||
for _, message := range messages {
|
||||
err := p.Produce(&confluent.Message{
|
||||
TopicPartition: confluent.TopicPartition{Topic: &topic, Partition: confluent.PartitionAny},
|
||||
Value: message,
|
||||
}, nil)
|
||||
if err != nil {
|
||||
return errors.Errorf("producing messages: %s", err)
|
||||
}
|
||||
p.Flush(3000)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
// deleteTopic deletes a kafka topic on the kafka server pointed to by
|
||||
// bootstrapServers
|
||||
func deleteTopic(topic string, bootstrapServers string) error {
|
||||
|
||||
// create admin client needed to create topic, defer closing
|
||||
ac, err := confluent.NewAdminClient(&confluent.ConfigMap{"bootstrap.servers": bootstrapServers})
|
||||
if err != nil {
|
||||
return errors.Errorf("creating admin client: %s", err)
|
||||
}
|
||||
defer ac.Close()
|
||||
|
||||
// delete a topic, catch any errors
|
||||
results, err := ac.DeleteTopics(context.Background(), []string{topic}, nil)
|
||||
if err != nil {
|
||||
return errors.Errorf("deleting topic: %s", err)
|
||||
}
|
||||
|
||||
for _, result := range results {
|
||||
if result.Error.Code() != confluent.ErrNoError {
|
||||
return errors.Errorf("fatal error deleting topic: %s", result.String())
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
}
|
||||
30
cli/kafka.go
30
cli/kafka.go
|
|
@ -1,35 +1,22 @@
|
|||
package cli
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/cli/batch"
|
||||
"github.com/featurebasedb/featurebase/v3/cli/kafka"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
func (cmd *Command) newKafkaRunner(cfgFile string) (*kafka.Runner, error) {
|
||||
// Read the kafka config file.
|
||||
v := viper.New()
|
||||
v.SetConfigFile(cfgFile)
|
||||
v.SetConfigType("toml")
|
||||
err := v.ReadInConfig()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error reading configuration file '%s': %v", cfgFile, err)
|
||||
}
|
||||
|
||||
cfg := kafka.Config{}
|
||||
if err := v.Unmarshal(&cfg); err != nil {
|
||||
return nil, errors.Wrap(err, "unmarshalling config")
|
||||
cfg, err := kafka.ConfigFromFile(cfgFile)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting config from file")
|
||||
}
|
||||
|
||||
if err := kafka.ValidateConfig(cfg); err != nil {
|
||||
return nil, errors.Wrap(err, "validating config")
|
||||
}
|
||||
|
||||
// Create a new config with defaults.
|
||||
|
||||
// Look up fields based on table provided in the config.
|
||||
wqr, err := cmd.executeQuery(newRawQuery("SHOW COLUMNS FROM " + cfg.Table))
|
||||
if err != nil {
|
||||
|
|
@ -57,14 +44,19 @@ func (cmd *Command) newKafkaRunner(cfgFile string) (*kafka.Runner, error) {
|
|||
return nil, errors.Wrap(err, "cleaning config")
|
||||
}
|
||||
|
||||
flds, err := kafka.ConfigToFields(cfg)
|
||||
flds, err := kafka.ConfigToFields(cfg, idkCfg.PrimaryKeys)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting fields from config")
|
||||
}
|
||||
|
||||
return kafka.NewRunner(
|
||||
kr := kafka.NewRunner(
|
||||
idkCfg,
|
||||
batch.NewSQLBatcher(cmd, flds),
|
||||
cmd.stderr,
|
||||
), nil
|
||||
)
|
||||
|
||||
// set pilosa host which is the only IDK config to come through
|
||||
// configuration flags rather than the kafka config file
|
||||
kr.Main.PilosaHosts = []string{cmd.host + ":" + cmd.port}
|
||||
return kr, nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,6 +8,14 @@ import (
|
|||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/idk"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
// Kafka message encoding types supported - changes here should be
|
||||
// propogated to the Config struct and ValidateConfig error message
|
||||
const (
|
||||
encodingTypeJSON = "json"
|
||||
encodingTypeAvro = "avro"
|
||||
)
|
||||
|
||||
// Config is the user-facing configuration for kafka support in the CLI. This is
|
||||
|
|
@ -23,6 +31,15 @@ type Config struct {
|
|||
|
||||
Table string `mapstructure:"table" help:"Destination table name."`
|
||||
Fields []Field `mapstructure:"fields"`
|
||||
|
||||
SchemaRegistryURL string `mapstructure:"schema-registry-url" help:"host and port of schema registry. Defaults to localhost:8081"`
|
||||
SchemaRegistryUsername string `mapstructure:"schema-registry-username" help:"authenticaion key provided by confluent for schema registry."`
|
||||
SchemaRegistryPassword string `mapstructure:"schema-registry-password" help:"authenticaion secret provided by confluent for schema registry."`
|
||||
|
||||
Encode string `mapstructure:"encode" help:"Encoding format (currently supported formats: avro, json)"`
|
||||
AllowMissingFields bool `mapstructure:"allow-missing-fields" help:"allow missing fields in messages from kafka"`
|
||||
MaxMessages int `mapstructure:"max-messages" help:"max messages read from kakfka"`
|
||||
ConfluentConfig string `mapstructure:"confluent-config" help:"path to JSON file mapping librdkafka consumer configurations to configuration values"`
|
||||
}
|
||||
|
||||
// Field is a user-facing configuration field.
|
||||
|
|
@ -44,72 +61,175 @@ type ConfigForIDK struct {
|
|||
BatchMaxStaleness time.Duration
|
||||
Timeout time.Duration
|
||||
|
||||
Table string
|
||||
IDField string
|
||||
Fields []idk.RawField
|
||||
Table string
|
||||
IDField string
|
||||
PrimaryKeys []string
|
||||
Fields []idk.RawField
|
||||
|
||||
SchemaRegistryURL string
|
||||
SchemaRegistryUsername string
|
||||
SchemaRegistryPassword string
|
||||
|
||||
Encode string
|
||||
AllowMissingFields bool
|
||||
MaxMessages int
|
||||
ConfluentConfig string
|
||||
}
|
||||
|
||||
// ValidateConfig validates the config is usable.
|
||||
// ConfigFromFile returns a Config struct based on a configuration file Default
|
||||
// values for Config struct are defined here
|
||||
func ConfigFromFile(cfgFile string) (cfg Config, err error) {
|
||||
// configure viper
|
||||
v := viper.New()
|
||||
v.SetConfigFile(cfgFile)
|
||||
v.SetConfigType("toml")
|
||||
|
||||
// set defaults
|
||||
v.SetDefault("hosts", []string{"localhost:9092"})
|
||||
v.SetDefault("group", "default-featurebase-group")
|
||||
v.SetDefault("batch-size", 1)
|
||||
v.SetDefault("batch-max-staleness", 5*time.Second)
|
||||
v.SetDefault("timeout", 5*time.Second)
|
||||
v.SetDefault("encode", encodingTypeJSON)
|
||||
|
||||
// Read the kafka config file.
|
||||
err = v.ReadInConfig()
|
||||
if err != nil {
|
||||
return cfg, fmt.Errorf("error reading configuration file '%s': %v", cfgFile, err)
|
||||
}
|
||||
|
||||
if err := v.Unmarshal(&cfg); err != nil {
|
||||
return cfg, errors.Wrap(err, "unmarshalling config")
|
||||
}
|
||||
|
||||
return
|
||||
|
||||
}
|
||||
|
||||
// ValidateConfig validates the config is usable. Note that different encoding
|
||||
// methods require different configurations
|
||||
func ValidateConfig(c Config) error {
|
||||
|
||||
// validate common
|
||||
if c.Table == "" {
|
||||
return errors.Errorf("table is required")
|
||||
} else if len(c.Topics) == 0 {
|
||||
}
|
||||
|
||||
if len(c.Topics) == 0 {
|
||||
return errors.Errorf("at least one topic is required")
|
||||
} else if len(c.Fields) > 0 {
|
||||
}
|
||||
|
||||
// validate on a by encoding basis
|
||||
switch c.Encode {
|
||||
case encodingTypeJSON:
|
||||
return validateConfigJSON(c)
|
||||
case encodingTypeAvro:
|
||||
return validateConfigAvro(c)
|
||||
default:
|
||||
return errors.Errorf("encode configuration value must be %s or %s: got %s", encodingTypeJSON, encodingTypeAvro, c.Encode)
|
||||
}
|
||||
}
|
||||
|
||||
func validateConfigJSON(c Config) error {
|
||||
|
||||
switch len(c.Fields) {
|
||||
case 0:
|
||||
// We only need to do these checks if any fields are specified at all.
|
||||
// If no fields are specified, that's ok because then we default to
|
||||
// using fields based off the existing table.
|
||||
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 c.Fields[i].Name == "" {
|
||||
return errors.Errorf("a name attribute (which isn't equal to \"\") should exist for all fields")
|
||||
}
|
||||
return nil
|
||||
case 1:
|
||||
return errors.Errorf("at least two fields are required (one should be a primary key)")
|
||||
default:
|
||||
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")
|
||||
if c.Fields[i].Name == "" {
|
||||
return errors.Errorf("a name attribute (which isn't equal to \"\") should exist for all fields")
|
||||
}
|
||||
if c.Fields[i].SourceType == "" {
|
||||
return errors.Errorf("a source-type attribute (which isn't equal to \"\") should exist for all fields")
|
||||
}
|
||||
}
|
||||
if found < 1 {
|
||||
return errors.Errorf("at least one primary key field is required")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
// Only primary key fields required
|
||||
func validateConfigAvro(c Config) error {
|
||||
|
||||
// for avro encoded messages, we just need to know what avro fields are
|
||||
// going to be used for the primary key. We'll check that there is at least
|
||||
// one field. For every field, we'll check that it has a name attribute and
|
||||
// has primary-key set.
|
||||
if len(c.Fields) < 1 {
|
||||
return errors.New("at least one field is required for avro encoded messages")
|
||||
}
|
||||
for i := range c.Fields {
|
||||
if !c.Fields[i].PrimaryKey {
|
||||
return errors.New("each field must be a primary key for avro encoded messages")
|
||||
}
|
||||
if c.Fields[i].Name == "" {
|
||||
return errors.New("a name attribute (which isn't equal to \"\") should exist for all fields")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
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,
|
||||
Hosts: c.Hosts,
|
||||
Group: c.Group,
|
||||
Topics: c.Topics,
|
||||
BatchSize: c.BatchSize,
|
||||
BatchMaxStaleness: c.BatchMaxStaleness,
|
||||
Timeout: c.Timeout,
|
||||
Table: c.Table,
|
||||
Encode: c.Encode,
|
||||
AllowMissingFields: c.AllowMissingFields,
|
||||
MaxMessages: c.MaxMessages,
|
||||
ConfluentConfig: c.ConfluentConfig,
|
||||
SchemaRegistryURL: c.SchemaRegistryURL,
|
||||
SchemaRegistryUsername: c.SchemaRegistryUsername,
|
||||
SchemaRegistryPassword: c.SchemaRegistryPassword,
|
||||
}
|
||||
|
||||
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 will be the same as c.Fields, but possibly enhanced.
|
||||
rawFields := make([]idk.RawField, 0, len(c.Fields))
|
||||
|
||||
var foundPK bool
|
||||
var stringKeys bool
|
||||
primaryKeys := []string{}
|
||||
for _, fld := range c.Fields {
|
||||
|
||||
// handle primary key
|
||||
if fld.PrimaryKey {
|
||||
out.IDField = fld.Name
|
||||
foundPK = true
|
||||
switch keyType := fld.SourceType; keyType {
|
||||
case dax.BaseTypeID:
|
||||
case dax.BaseTypeString, dax.BaseTypeInt:
|
||||
stringKeys = true
|
||||
default:
|
||||
// IDK can handle other field types as primary keys but limiting
|
||||
// here to the ones above for now. ID and string are the ones
|
||||
// that make sense and existing users also use int fields so I'm
|
||||
// including that as well.
|
||||
return out, errors.Errorf("primary-key fields must be \"id\", \"string\", or \"int\": got field %s which is type %s", fld.Name, keyType)
|
||||
}
|
||||
primaryKeys = append(primaryKeys, fld.Name)
|
||||
}
|
||||
|
||||
typ, quals, err := dax.SplitFieldType(fld.SourceType)
|
||||
|
|
@ -149,8 +269,16 @@ func ConvertConfig(c Config) (ConfigForIDK, error) {
|
|||
|
||||
rawFields = append(rawFields, rawFld)
|
||||
}
|
||||
if !foundPK {
|
||||
|
||||
// Should have at least one primary key. If there is more than one primary key OR
|
||||
// using a string field as the key then use string keys (i.e. table will be keyed)
|
||||
// Else, use ids (i.e. table will not be keyed)
|
||||
if len(primaryKeys) < 1 {
|
||||
return out, errors.New("primary-key not found in fields")
|
||||
} else if stringKeys || len(primaryKeys) > 1 {
|
||||
out.PrimaryKeys = primaryKeys
|
||||
} else {
|
||||
out.IDField = primaryKeys[0]
|
||||
}
|
||||
|
||||
out.Fields = rawFields
|
||||
|
|
@ -160,13 +288,22 @@ func ConvertConfig(c Config) (ConfigForIDK, error) {
|
|||
|
||||
// ConfigToFields returns a list of *dax.Field based on the IDField and Fields
|
||||
// in the Config.
|
||||
func ConfigToFields(c Config) ([]*dax.Field, error) {
|
||||
func ConfigToFields(c Config, primaryKeys []string) ([]*dax.Field, error) {
|
||||
// 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 avro, let the SchemaManager and IDK handle fields
|
||||
if c.Encode == encodingTypeAvro {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
for _, fld := range c.Fields {
|
||||
if fld.PrimaryKey {
|
||||
// When we have a single primary key, don't also store that value as a
|
||||
// field in FeatureBase. However, when we have more than one primary
|
||||
// key, store all the values used for the compound key as fields in
|
||||
// FeatureBase.
|
||||
if fld.PrimaryKey && len(primaryKeys) < 2 {
|
||||
continue
|
||||
}
|
||||
typ, quals, err := dax.SplitFieldType(fld.SourceType)
|
||||
|
|
|
|||
47
cli/kafka/config_test.go
Normal file
47
cli/kafka/config_test.go
Normal file
|
|
@ -0,0 +1,47 @@
|
|||
package kafka
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type configFromFileTest struct {
|
||||
configFilePath string
|
||||
expectedStruct string
|
||||
}
|
||||
|
||||
func TestConfigFromFile(t *testing.T) {
|
||||
tests := []configFromFileTest{
|
||||
{ // confirm json config struct fields are being set
|
||||
configFilePath: "./testdata/config/config01.toml",
|
||||
expectedStruct: `{"Hosts":["kafka:9090"],"Group":"testGroup","Topics":["testTopic"],"BatchSize":35,"BatchMaxStaleness":25000000000,"Timeout":16000000000,"Table":"testTable","Fields":[{"Name":"id","SourceType":"id","SourcePath":null,"PrimaryKey":true},{"Name":"name","SourceType":"string","SourcePath":["test","path"],"PrimaryKey":false},{"Name":"age","SourceType":"int","SourcePath":null,"PrimaryKey":false},{"Name":"hobbies","SourceType":"stringset","SourcePath":null,"PrimaryKey":false}],"SchemaRegistryURL":"","SchemaRegistryUsername":"","SchemaRegistryPassword":"","Encode":"json","AllowMissingFields":false,"MaxMessages":0,"ConfluentConfig":""}`,
|
||||
},
|
||||
{ // confirm defaults are being set
|
||||
configFilePath: "./testdata/config/config00.toml",
|
||||
expectedStruct: `{"Hosts":["localhost:9092"],"Group":"default-featurebase-group","Topics":["topic00"],"BatchSize":1,"BatchMaxStaleness":5000000000,"Timeout":5000000000,"Table":"table00","Fields":[{"Name":"id","SourceType":"string","SourcePath":null,"PrimaryKey":true},{"Name":"name","SourceType":"string","SourcePath":null,"PrimaryKey":false},{"Name":"age","SourceType":"int","SourcePath":null,"PrimaryKey":false},{"Name":"hobbies","SourceType":"stringset","SourcePath":null,"PrimaryKey":false}],"SchemaRegistryURL":"","SchemaRegistryUsername":"","SchemaRegistryPassword":"","Encode":"json","AllowMissingFields":false,"MaxMessages":0,"ConfluentConfig":""}`,
|
||||
},
|
||||
{ // confirm avro config struct field are being set
|
||||
configFilePath: "./testdata/config/config02.toml",
|
||||
expectedStruct: `{"Hosts":["kafka:9090"],"Group":"testGroup","Topics":["testTopic"],"BatchSize":35,"BatchMaxStaleness":25000000000,"Timeout":16000000000,"Table":"testTable","Fields":[{"Name":"id","SourceType":"id","SourcePath":null,"PrimaryKey":true},{"Name":"name","SourceType":"string","SourcePath":["test","path"],"PrimaryKey":false},{"Name":"age","SourceType":"int","SourcePath":null,"PrimaryKey":false},{"Name":"hobbies","SourceType":"stringset","SourcePath":null,"PrimaryKey":false}],"SchemaRegistryURL":"","SchemaRegistryUsername":"","SchemaRegistryPassword":"","Encode":"avro","AllowMissingFields":false,"MaxMessages":0,"ConfluentConfig":""}`,
|
||||
},
|
||||
{ // confirm avro config struct field are being set
|
||||
configFilePath: "./testdata/config/config03.toml",
|
||||
expectedStruct: `{"Hosts":["kafka:9090"],"Group":"testGroup","Topics":["testTopic"],"BatchSize":35,"BatchMaxStaleness":25000000000,"Timeout":16000000000,"Table":"testTable","Fields":[{"Name":"id","SourceType":"id","SourcePath":null,"PrimaryKey":true},{"Name":"name","SourceType":"string","SourcePath":["test","path"],"PrimaryKey":false},{"Name":"age","SourceType":"int","SourcePath":null,"PrimaryKey":false},{"Name":"hobbies","SourceType":"stringset","SourcePath":null,"PrimaryKey":false}],"SchemaRegistryURL":"","SchemaRegistryUsername":"","SchemaRegistryPassword":"","Encode":"avro","AllowMissingFields":true,"MaxMessages":100,"ConfluentConfig":"./test/confluent/config.json"}`,
|
||||
},
|
||||
}
|
||||
|
||||
for _, test := range tests {
|
||||
cfg, err := ConfigFromFile(test.configFilePath)
|
||||
if err != nil {
|
||||
t.Fatalf("getting config struct from file: %s", err)
|
||||
}
|
||||
json, err := json.Marshal(cfg)
|
||||
if err != nil {
|
||||
t.Fatalf("marshaling config struct: %s", err)
|
||||
}
|
||||
require.JSONEq(t, test.expectedStruct, string(json))
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -7,7 +7,9 @@ import (
|
|||
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"
|
||||
"github.com/featurebasedb/featurebase/v3/idk/common"
|
||||
"github.com/featurebasedb/featurebase/v3/idk/kafka"
|
||||
"github.com/featurebasedb/featurebase/v3/idk/kafka_sasl"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
)
|
||||
|
||||
|
|
@ -15,53 +17,105 @@ import (
|
|||
// idk.kafka_static.Main in that it embeds idk.Main and contains additional
|
||||
// functionality specific to its use case.
|
||||
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."`
|
||||
idk.Main `flag:"!embed"`
|
||||
idk.ConfluentCommand `flag:"!embed"`
|
||||
Hosts []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."`
|
||||
}
|
||||
|
||||
// Configure and return *runner with common elements of all kafka runners
|
||||
func NewRunner(cfg ConfigForIDK, batcher fbbatch.Batcher, logWriter io.Writer) *Runner {
|
||||
idkMain := idk.NewMain()
|
||||
idkMain.IDField = cfg.IDField
|
||||
idkMain.PrimaryKeyFields = cfg.PrimaryKeys
|
||||
idkMain.Index = cfg.Table
|
||||
idkMain.Batcher = batcher
|
||||
idkMain.BatchSize = cfg.BatchSize
|
||||
idkMain.BatchMaxStaleness = cfg.BatchMaxStaleness
|
||||
idkMain.MaxMsgs = uint64(cfg.MaxMessages)
|
||||
idkMain.SetBasic()
|
||||
idkMain.SetLog(logger.NewStandardLogger(logWriter))
|
||||
idkMain.OffsetMode = true
|
||||
idkMain.Namespace = "sql_kafka_runner"
|
||||
idkMain.Pprof = "" // don't initialize pprof until we actually use it in tests
|
||||
|
||||
kr := &Runner{
|
||||
Main: *idkMain,
|
||||
KafkaHosts: cfg.Hosts,
|
||||
Group: cfg.Group,
|
||||
Topics: cfg.Topics,
|
||||
Header: cfg.Fields,
|
||||
Timeout: cfg.Timeout,
|
||||
Main: *idkMain,
|
||||
Hosts: 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()
|
||||
kr.KafkaBootstrapServers = cfg.Hosts
|
||||
kr.SchemaRegistryURL = cfg.SchemaRegistryURL
|
||||
kr.SchemaRegistryUsername = cfg.SchemaRegistryUsername
|
||||
kr.SchemaRegistryPassword = cfg.SchemaRegistryPassword
|
||||
|
||||
// NewSource should be set based on the encoding of the source (e.g. JSON, Avro)
|
||||
switch cfg.Encode {
|
||||
case encodingTypeAvro:
|
||||
kr.GetAvroNewSource(cfg)
|
||||
default:
|
||||
kr.GetJSONNewSource(cfg)
|
||||
}
|
||||
|
||||
return kr
|
||||
}
|
||||
|
||||
func (r *Runner) GetJSONNewSource(cfg ConfigForIDK) {
|
||||
|
||||
r.NewSource = func() (idk.Source, error) {
|
||||
source := kafka_sasl.NewSource()
|
||||
//source.KafkaBootstrapServers = r.Hosts
|
||||
source.Group = r.Group
|
||||
source.Topics = r.Topics
|
||||
source.Log = r.Main.Log()
|
||||
source.Timeout = r.Timeout
|
||||
source.HeaderFields = r.Header
|
||||
source.AllowMissingFields = cfg.AllowMissingFields
|
||||
source.KafkaConfiguration = cfg.ConfluentConfig
|
||||
confluentCfg, err := common.SetupConfluent(&r.ConfluentCommand)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
source.ConfigMap = confluentCfg
|
||||
|
||||
err = source.Open()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "opening source")
|
||||
}
|
||||
return source, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Runner) GetAvroNewSource(cfg ConfigForIDK) {
|
||||
|
||||
r.NewSource = func() (idk.Source, error) {
|
||||
source := kafka.NewSource()
|
||||
source.KafkaBootstrapServers = r.Hosts
|
||||
source.SchemaRegistryURL = cfg.SchemaRegistryURL
|
||||
source.SchemaRegistryUsername = cfg.SchemaRegistryUsername
|
||||
source.SchemaRegistryPassword = cfg.SchemaRegistryPassword
|
||||
source.Group = r.Group
|
||||
source.Topics = r.Topics
|
||||
source.Log = r.Main.Log()
|
||||
source.Timeout = r.Timeout
|
||||
source.KafkaConfiguration = cfg.ConfluentConfig
|
||||
confluentCfg, err := common.SetupConfluent(&r.ConfluentCommand)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
source.ConfigMap = confluentCfg
|
||||
|
||||
err = source.Open()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "opening source")
|
||||
}
|
||||
return source, nil
|
||||
}
|
||||
return kr
|
||||
}
|
||||
|
|
|
|||
20
cli/kafka/testdata/config/config00.toml
vendored
Normal file
20
cli/kafka/testdata/config/config00.toml
vendored
Normal file
|
|
@ -0,0 +1,20 @@
|
|||
## confirm defaults are being set correctly
|
||||
topics = "topic00"
|
||||
table = "table00"
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
26
cli/kafka/testdata/config/config01.toml
vendored
Normal file
26
cli/kafka/testdata/config/config01.toml
vendored
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
hosts = ["kafka:9090"]
|
||||
group = "testGroup"
|
||||
topics = "testTopic"
|
||||
table = "testTable"
|
||||
batch-size = 35
|
||||
batch-max-staleness = "25s"
|
||||
timeout = "16s"
|
||||
encode = "json"
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
source-path = ["test", "path"]
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
27
cli/kafka/testdata/config/config02.toml
vendored
Normal file
27
cli/kafka/testdata/config/config02.toml
vendored
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
hosts = ["kafka:9090"]
|
||||
group = "testGroup"
|
||||
topics = "testTopic"
|
||||
table = "testTable"
|
||||
batch-size = 35
|
||||
batch-max-staleness = "25s"
|
||||
timeout = "16s"
|
||||
encode = "avro"
|
||||
schemaRegistryHost = "localhost:8081"
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
source-path = ["test", "path"]
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
30
cli/kafka/testdata/config/config03.toml
vendored
Normal file
30
cli/kafka/testdata/config/config03.toml
vendored
Normal file
|
|
@ -0,0 +1,30 @@
|
|||
hosts = ["kafka:9090"]
|
||||
group = "testGroup"
|
||||
topics = "testTopic"
|
||||
table = "testTable"
|
||||
batch-size = 35
|
||||
batch-max-staleness = "25s"
|
||||
timeout = "16s"
|
||||
encode = "avro"
|
||||
schemaRegistryHost = "localhost:8081"
|
||||
max-messages = 100
|
||||
allow-missing-fields = true
|
||||
confluent-config = "./test/confluent/config.json"
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
source-path = ["test", "path"]
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
26
cli/kafka/testdata/runner/config/config00.toml
vendored
Normal file
26
cli/kafka/testdata/runner/config/config00.toml
vendored
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic00"
|
||||
table = "table00"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "json"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
26
cli/kafka/testdata/runner/config/config01.toml
vendored
Normal file
26
cli/kafka/testdata/runner/config/config01.toml
vendored
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic01"
|
||||
table = "table01"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "json"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
27
cli/kafka/testdata/runner/config/config02.toml
vendored
Normal file
27
cli/kafka/testdata/runner/config/config02.toml
vendored
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic02"
|
||||
table = "table02"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "json"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
28
cli/kafka/testdata/runner/config/config03.toml
vendored
Normal file
28
cli/kafka/testdata/runner/config/config03.toml
vendored
Normal file
|
|
@ -0,0 +1,28 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic03"
|
||||
table = "table03"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "json"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
15
cli/kafka/testdata/runner/config/config04.toml
vendored
Normal file
15
cli/kafka/testdata/runner/config/config04.toml
vendored
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic04"
|
||||
table = "table04"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "avro"
|
||||
schema-registry-url = "SCHEMA_REGISTRY_SERVICE"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "pk0"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
27
cli/kafka/testdata/runner/config/config05.toml
vendored
Normal file
27
cli/kafka/testdata/runner/config/config05.toml
vendored
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic05"
|
||||
table = "table05"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "json"
|
||||
allow-missing-fields = true
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "id"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "name"
|
||||
source-type = "string"
|
||||
|
||||
[[fields]]
|
||||
name = "age"
|
||||
source-type = "int"
|
||||
|
||||
[[fields]]
|
||||
name = "hobbies"
|
||||
source-type = "stringset"
|
||||
15
cli/kafka/testdata/runner/config/config06.toml
vendored
Normal file
15
cli/kafka/testdata/runner/config/config06.toml
vendored
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic06"
|
||||
table = "table06"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "avro"
|
||||
schema-registry-url = "SCHEMA_REGISTRY_SERVICE"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "pk"
|
||||
source-type = "id"
|
||||
primary-key = true
|
||||
26
cli/kafka/testdata/runner/config/config07.toml
vendored
Normal file
26
cli/kafka/testdata/runner/config/config07.toml
vendored
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
hosts = ["KAFKA_SERVICE"]
|
||||
group = "grp"
|
||||
topics = "topic07"
|
||||
table = "table07"
|
||||
batch-size = 1
|
||||
batch-max-staleness = "5s"
|
||||
timeout = "5s"
|
||||
encode = "avro"
|
||||
schema-registry-url = "SCHEMA_REGISTRY_SERVICE"
|
||||
max-messages = MAX_MESSAGES
|
||||
|
||||
[[fields]]
|
||||
name = "pk0"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
|
||||
[[fields]]
|
||||
name = "pk1"
|
||||
source-type = "string"
|
||||
primary-key = true
|
||||
|
||||
[[fields]]
|
||||
name = "int_int"
|
||||
source-type = "int"
|
||||
primary-key = true
|
||||
6
cli/kafka/testdata/runner/data/data00.json
vendored
Normal file
6
cli/kafka/testdata/runner/data/data00.json
vendored
Normal file
|
|
@ -0,0 +1,6 @@
|
|||
{"id": 1, "name": "a", "age": 20, "hobbies": "hob1,hob2"}
|
||||
{"id": 2, "age": 21, "hobbies": "hob3,hob2", "name": "b"}
|
||||
{"id": 3, "name": "c", "age": 22, "hobbies": "hob3,hob4"}
|
||||
{"id": 4, "name": "d", "age": 23, "hobbies": "hob4,hob5"}
|
||||
{"id": 5, "name": "e", "age": 24, "hobbies": "hob5,hob6"}
|
||||
{"id": 6, "name": "f", "age": 26, "hobbies": "hob6,hob7"}
|
||||
10
cli/kafka/testdata/runner/data/data01.json
vendored
Normal file
10
cli/kafka/testdata/runner/data/data01.json
vendored
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
{"pk0": "9z4aw", "pk1": "5ptDx", "pk2": "CKs1F", "stringset_string": {"string": "7EYSp"}, "string_string": {"string": "uirDR"}, "stringset_bytes": {"bytes": "gL2Hg"}, "string_bytes": {"bytes": "BmvHF"}, "stringset_stringarray": {"array": ["vbbuf", "VQs7y", "9z4aw", "h1iqc", "aQQxr"]}, "stringset_bytesarray": {"array": ["u2Yr4", "tvNOB", "iYeOV", "ZgkOB", "RPGAm"]}, "idset_long": {"long": 647}, "id_long": {"long": 792}, "idset_int": {"int": 898}, "id_int": {"int": 63}, "idset_longarray": {"array": [442, 167, 230, 344, 733]}, "idset_intarray": {"array": [442, 614, 394, 284, 344]}, "int_long": {"long": 584}, "int_int": {"int": 344}, "decimal_bytes": {"bytes": "\u0001\u00cb"}, "decimal_float": {"float": 3.23}, "decimal_double": {"double": 0.95}, "dateint_bytes_ts": {"bytes": "2023-02-16 07:53:59"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-02-16 07:53:59"}, "timestamp_bytes_int": {"bytes": "1676555639"}}
|
||||
{"pk0": "ASSAw", "pk1": "kauLy", "pk2": "oxjI0", "stringset_string": {"string": "iYeOV"}, "string_string": {"string": "LBTEU"}, "stringset_bytes": {"bytes": "5ptDx"}, "string_bytes": {"bytes": "EyQoi"}, "stringset_stringarray": {"array": ["iYeOV", "XzEHj", "rrkYB", "v31XN", "uirDR"]}, "stringset_bytesarray": {"array": ["X9jWC", "x5z8P", "PYE8V", "PYE8V", "vTwn4"]}, "idset_long": {"long": 484}, "id_long": {"long": 23}, "idset_int": {"int": 322}, "id_int": {"int": 320}, "idset_longarray": {"array": [792, 809, 168, 399, 639]}, "idset_intarray": {"array": [606, 293, 23, 358, 821]}, "int_long": {"long": 533}, "int_int": {"int": 884}, "decimal_bytes": {"bytes": "\u0001\u0029"}, "decimal_float": {"float": 4.32}, "decimal_double": {"double": 4.97}, "dateint_bytes_ts": {"bytes": "2023-02-22 14:32:23"}, "bool_bool": {"boolean": true}, "timestamp_bytes_ts": {"bytes": "2023-02-22 14:32:23"}, "timestamp_bytes_int": {"bytes": "1677097943"}}
|
||||
{"pk0": "BmvHF", "pk1": "798ka", "pk2": "6TKzc", "stringset_string": {"string": "X9jWC"}, "string_string": {"string": "I1gXJ"}, "stringset_bytes": {"bytes": "gjWEI"}, "string_bytes": {"bytes": "thuky"}, "stringset_stringarray": {"array": ["tyP3m", "5ptDx", "TLaUE", "EyQoi", "Chgzr"]}, "stringset_bytesarray": {"array": ["VQs7y", "h1iqc", "F0uC4", "d0U7s", "byHh9"]}, "idset_long": {"long": 166}, "id_long": {"long": 320}, "idset_int": {"int": 232}, "id_int": {"int": 286}, "idset_longarray": {"array": [890, 975, 284, 289, 388]}, "idset_intarray": {"array": [865, 931, 614, 884, 322]}, "int_long": {"long": 857}, "int_int": {"int": 879}, "decimal_bytes": {"bytes": "\u0000\u00e6"}, "decimal_float": {"float": 2.84}, "decimal_double": {"double": 2.91}, "dateint_bytes_ts": {"bytes": "2023-01-31 11:11:30"}, "bool_bool": {"boolean": true }, "timestamp_bytes_ts": {"bytes": "2023-01-31 11:11:30"}, "timestamp_bytes_int": {"bytes": "1675185090"}}
|
||||
{"pk0": "tElMR", "pk1": "ARlcJ", "pk2": "n9HUP", "stringset_string": {"string": "58KIR"}, "string_string": {"string": "FW39I"}, "stringset_bytes": {"bytes": "PNB4s"}, "string_bytes": {"bytes": "FW39I"}, "stringset_stringarray": {"array": ["X9jWC", "58KIR", "X9jWC", "6TKzc", "8MGwy"]}, "stringset_bytesarray": {"array": ["vhisL", "BmvHF", "eofzb", "TLaUE", "PNB4s"]}, "idset_long": {"long": 289}, "id_long": {"long": 695}, "idset_int": {"int": 791}, "id_int": {"int": 821}, "idset_longarray": {"array": [2, 680, 958, 289, 389]}, "idset_intarray": {"array": [606, 890, 387, 102, 220]}, "int_long": {"long": 289}, "int_int": {"int": 2}, "decimal_bytes": {"bytes": "\u005f"}, "decimal_float": {"float": 2.65}, "decimal_double": {"double": 2.19}, "dateint_bytes_ts": {"bytes": "2023-02-20 17:04:21"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-02-20 17:04:21"}, "timestamp_bytes_int": {"bytes": "1676934261"}}
|
||||
{"pk0": "RKE3c", "pk1": "6TKzc", "pk2": "RKE3c", "stringset_string": {"string": "dF6kx"}, "string_string": {"string": "TLaUE"}, "stringset_bytes": {"bytes": "dxKKn"}, "string_bytes": {"bytes": "YdwQY"}, "stringset_stringarray": {"array": ["5HIn2", "wNZ7o", "KdTtE", "x5z8P", "nVQrd"]}, "stringset_bytesarray": {"array": ["5HIn2", "tyP3m", "I6NST", "gjWEI", "Qylqq"]}, "idset_long": {"long": 39}, "id_long": {"long": 665}, "idset_int": {"int": 113}, "id_int": {"int": 681}, "idset_longarray": {"array": [857, 63, 172, 220, 358]}, "idset_intarray": {"array": [220, 731, 647, 778, 665]}, "int_long": {"long": 582}, "int_int": {"int": 690}, "decimal_bytes": {"bytes": "\u0001\u00d1"}, "decimal_float": {"float": 1.36}, "decimal_double": {"double": 2.72}, "dateint_bytes_ts": {"bytes": "2023-02-16 20:09:13"}, "bool_bool": {"boolean": true }, "timestamp_bytes_ts": {"bytes": "2023-02-16 20:09:13"}, "timestamp_bytes_int": {"bytes": "1676599753"}}
|
||||
{"pk0": "yg8hY", "pk1": "tvNOB", "pk2": "byHh9", "stringset_string": {"string": "911oj"}, "string_string": {"string": "5HIn2"}, "stringset_bytes": {"bytes": "u2Yr4"}, "string_bytes": {"bytes": "qK5TE"}, "stringset_stringarray": {"array": ["nVQrd", "fjQK2", "m5d59", "dxKKn", "d0U7s"]}, "stringset_bytesarray": {"array": ["wNZ7o", "OKNV2", "F0uC4", "VBcyJ", "KMZnH"]}, "idset_long": {"long": 839}, "id_long": {"long": 809}, "idset_int": {"int": 533}, "id_int": {"int": 168}, "idset_longarray": {"array": [582, 629, 680, 63, 690]}, "idset_intarray": {"array": [969, 175, 172, 257, 115]}, "int_long": {"long": 433}, "int_int": {"int": 680}, "decimal_bytes": {"bytes": "\u0001\u00d6"}, "decimal_float": {"float": 1.27}, "decimal_double": {"double": 3.23}, "dateint_bytes_ts": {"bytes": "2023-01-30 06:56:05"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-01-30 06:56:05"}, "timestamp_bytes_int": {"bytes": "1675083365"}}
|
||||
{"pk0": "6TKzc", "pk1": "YKLk9", "pk2": "h1iqc", "stringset_string": {"string": "eofzb"}, "string_string": {"string": "n9HUP"}, "stringset_bytes": {"bytes": "t5f7R"}, "string_bytes": {"bytes": "5HIn2"}, "stringset_stringarray": {"array": ["yg8hY", "xE5jX", "C6xxn", "BmvHF", "PYE8V"]}, "stringset_bytesarray": {"array": ["6TKzc", "vK0WD", "xE5jX", "jVVfZ", "pjxqm"]}, "idset_long": {"long": 72}, "id_long": {"long": 387}, "idset_int": {"int": 676}, "id_int": {"int": 797}, "idset_longarray": {"array": [399, 322, 975, 730, 969]}, "idset_intarray": {"array": [958, 242, 778, 289, 797]}, "int_long": {"long": 430}, "int_int": {"int": 23}, "decimal_bytes": {"bytes": "\u0001\u00d6"}, "decimal_float": {"float": 3.35}, "decimal_double": {"double": 0.78}, "dateint_bytes_ts": {"bytes": "2023-02-23 05:04:34"}, "bool_bool": {"boolean": true }, "timestamp_bytes_ts": {"bytes": "2023-02-23 05:04:34"}, "timestamp_bytes_int": {"bytes": "1677150274"}}
|
||||
{"pk0": "h1iqc", "pk1": "5ptDx", "pk2": "iYeOV", "stringset_string": {"string": "ASSAw"}, "string_string": {"string": "58KIR"}, "stringset_bytes": {"bytes": "eNKWF"}, "string_bytes": {"bytes": "x5z8P"}, "stringset_stringarray": {"array": ["pjxqm", "6TKzc", "ZgkOB", "eofzb", "RKE3c"]}, "stringset_bytesarray": {"array": ["BmvHF", "Qylqq", "5HIn2", "7EYSp", "yTeUQ"]}, "idset_long": {"long": 255}, "id_long": {"long": 647}, "idset_int": {"int": 821}, "id_int": {"int": 389}, "idset_longarray": {"array": [606, 110, 320, 63, 344]}, "idset_intarray": {"array": [29, 289, 388, 257, 606]}, "int_long": {"long": 110}, "int_int": {"int": 148}, "decimal_bytes": {"bytes": "\u0001\u001c"}, "decimal_float": {"float": 0.4}, "decimal_double": {"double": 4.41}, "dateint_bytes_ts": {"bytes": "2023-02-19 08:52:56"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-02-19 08:52:56"}, "timestamp_bytes_int": {"bytes": "1676818376"}}
|
||||
{"pk0": "DY2Ui", "pk1": "kUbdU", "pk2": "pjxqm", "stringset_string": {"string": "tyP3m"}, "string_string": {"string": "8MGwy"}, "stringset_bytes": {"bytes": "DDLN5"}, "string_bytes": {"bytes": "vTwn4"}, "stringset_stringarray": {"array": ["XzEHj", "8MGwy", "gjWEI", "xE5jX", "v31XN"]}, "stringset_bytesarray": {"array": ["d0U7s", "u2Yr4", "d0U7s", "sDdtS", "y2Y7b"]}, "idset_long": {"long": 984}, "id_long": {"long": 430}, "idset_int": {"int": 931}, "id_int": {"int": 297}, "idset_longarray": {"array": [975, 733, 113, 751, 772]}, "idset_intarray": {"array": [297, 72, 694, 898, 384]}, "int_long": {"long": 63}, "int_int": {"int": 388}, "decimal_bytes": {"bytes": "\u007e"}, "decimal_float": {"float": 0.83}, "decimal_double": {"double": 4.23}, "dateint_bytes_ts": {"bytes": "2023-02-12 18:37:16"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-02-12 18:37:16"}, "timestamp_bytes_int": {"bytes": "1676248636"}}
|
||||
{"pk0": "u2Yr4", "pk1": "sHaUv", "pk2": "x5z8P", "stringset_string": {"string": "7EYSp"}, "string_string": {"string": "ZgkOB"}, "stringset_bytes": {"bytes": "qK5TE"}, "string_bytes": {"bytes": "6iGIm"}, "stringset_stringarray": {"array": ["u2Yr4", "PYE8V", "VBcyJ", "Chgzr", "DY2Ui"]}, "stringset_bytesarray": {"array": ["YdwQY", "kUbdU", "aQQxr", "KdTtE", "MVNow"]}, "idset_long": {"long": 148}, "id_long": {"long": 115}, "idset_int": {"int": 890}, "id_int": {"int": 39}, "idset_longarray": {"array": [839, 63, 148, 984, 958]}, "idset_intarray": {"array": [731, 13, 167, 772, 629]}, "int_long": {"long": 13}, "int_int": {"int": 969}, "decimal_bytes": {"bytes": "\u0000\u009a"}, "decimal_float": {"float": 2.93}, "decimal_double": {"double": 2.29}, "dateint_bytes_ts": {"bytes": "2023-02-03 16:19:37"}, "bool_bool": {"boolean": false}, "timestamp_bytes_ts": {"bytes": "2023-02-03 16:19:37"}, "timestamp_bytes_int": {"bytes": "1675462777"}}
|
||||
6
cli/kafka/testdata/runner/data/data02.json
vendored
Normal file
6
cli/kafka/testdata/runner/data/data02.json
vendored
Normal file
|
|
@ -0,0 +1,6 @@
|
|||
{"id": 1, "name": "a", "age": 20, "hobbies": "hob1,hob2"}
|
||||
{"id": 2, "age": 21, "hobbies": "hob3,hob2", "name": "b"}
|
||||
{"id": 3, "age": 22, "hobbies": "hob3,hob4"}
|
||||
{"id": 4, "name": "d", "hobbies": "hob4,hob5"}
|
||||
{"id": 5, "name": "e", "age": 24}
|
||||
{"id": 6, "name": "f", "age": 26, "hobbies": "hob6,hob7"}
|
||||
10
cli/kafka/testdata/runner/data/data03.json
vendored
Normal file
10
cli/kafka/testdata/runner/data/data03.json
vendored
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
{"pk": 0, "string_string": {"string": "uirDR"}, "stringset_stringarray": {"array": ["vbbuf", "VQs7y", "9z4aw", "h1iqc", "aQQxr"]}, "idset_longarray": {"array": [442, 167, 230]}, "int_int": {"int": 344}, "bool_bool": {"boolean": false}}
|
||||
{"pk": 2, "string_string": {"string": "LBTEU"}, "stringset_stringarray": {"array": ["iYeOV", "XzEHj", "rrkYB", "v31XN", "uirDR"]}, "idset_longarray": {"array": [792, 809, 168]}, "int_int": {"int": 884}, "bool_bool": {"boolean": true }}
|
||||
{"pk": 4, "string_string": {"string": "I1gXJ"}, "stringset_stringarray": {"array": ["tyP3m", "5ptDx", "TLaUE", "EyQoi", "Chgzr"]}, "idset_longarray": {"array": [890, 975, 284]}, "int_int": {"int": 879}, "bool_bool": {"boolean": true }}
|
||||
{"pk": 6, "string_string": {"string": "FW39I"}, "stringset_stringarray": {"array": ["X9jWC", "58KIR", "X9jWC", "6TKzc", "8MGwy"]}, "idset_longarray": {"array": [201, 680, 958]}, "int_int": {"int": 212}, "bool_bool": {"boolean": false}}
|
||||
{"pk": 8, "string_string": {"string": "TLaUE"}, "stringset_stringarray": {"array": ["5HIn2", "wNZ7o", "KdTtE", "x5z8P", "nVQrd"]}, "idset_longarray": {"array": [857, 630, 172]}, "int_int": {"int": 690}, "bool_bool": {"boolean": true }}
|
||||
{"pk": 10, "string_string": {"string": "5HIn2"}, "stringset_stringarray": {"array": ["nVQrd", "fjQK2", "m5d59", "dxKKn", "d0U7s"]}, "idset_longarray": {"array": [582, 629, 680]}, "int_int": {"int": 680}, "bool_bool": {"boolean": false}}
|
||||
{"pk": 12, "string_string": {"string": "n9HUP"}, "stringset_stringarray": {"array": ["yg8hY", "xE5jX", "C6xxn", "BmvHF", "PYE8V"]}, "idset_longarray": {"array": [399, 322, 975]}, "int_int": {"int": 230}, "bool_bool": {"boolean": true }}
|
||||
{"pk": 14, "string_string": {"string": "58KIR"}, "stringset_stringarray": {"array": ["pjxqm", "6TKzc", "ZgkOB", "eofzb", "RKE3c"]}, "idset_longarray": {"array": [606, 110, 320]}, "int_int": {"int": 148}, "bool_bool": {"boolean": false}}
|
||||
{"pk": 16, "string_string": {"string": "8MGwy"}, "stringset_stringarray": {"array": ["XzEHj", "8MGwy", "gjWEI", "xE5jX", "v31XN"]}, "idset_longarray": {"array": [975, 733, 113]}, "int_int": {"int": 388}, "bool_bool": {"boolean": false}}
|
||||
{"pk": 18, "string_string": {"string": "ZgkOB"}, "stringset_stringarray": {"array": ["u2Yr4", "PYE8V", "VBcyJ", "Chgzr", "DY2Ui"]}, "idset_longarray": {"array": [839, 635, 148]}, "int_int": {"int": 969}, "bool_bool": {"boolean": false}}
|
||||
32
cli/kafka/testdata/runner/schema/schema01.json
vendored
Normal file
32
cli/kafka/testdata/runner/schema/schema01.json
vendored
Normal file
|
|
@ -0,0 +1,32 @@
|
|||
{
|
||||
"namespace": "org.test",
|
||||
"type": "record",
|
||||
"name": "all_type_schema",
|
||||
"doc": "All supported avro types and property variations",
|
||||
"fields": [
|
||||
{"name": "pk0", "type": "string"},
|
||||
{"name": "pk1", "type": "string"},
|
||||
{"name": "pk2", "type": "string"},
|
||||
{"name": "stringset_string", "type": ["string", "null"], "mutex": false },
|
||||
{"name": "string_string", "type": ["string", "null"], "mutex": true },
|
||||
{"name": "stringset_bytes", "type": ["bytes", "null"], "mutex": false},
|
||||
{"name": "string_bytes", "type": ["bytes", "null"] , "mutex": true },
|
||||
{"name": "stringset_stringarray", "type": [{"type": "array", "items": "string"}, "null"]},
|
||||
{"name": "stringset_bytesarray", "type": [{"type": "array", "items": "string"}, "null"]},
|
||||
{"name": "idset_long", "type": ["long", "null"], "mutex": false, "fieldType": "id"},
|
||||
{"name": "id_long", "type": ["long", "null"], "mutex": true, "fieldType": "id"},
|
||||
{"name": "idset_int", "type": ["int", "null"], "mutex": false, "fieldType": "id"},
|
||||
{"name": "id_int", "type": ["int", "null"], "mutex": true, "fieldType": "id"},
|
||||
{"name": "idset_longarray", "type": [{"type": "array", "items": "long"}, "null"], "fieldType": "id"},
|
||||
{"name": "idset_intarray", "type": [{"type": "array", "items": "int"}, "null"]},
|
||||
{"name": "int_long", "type": ["long", "null"], "fieldType": "int"},
|
||||
{"name": "int_int", "type": ["int", "null"], "fieldType": "int"},
|
||||
{"name": "decimal_bytes", "type": ["bytes", "null"], "fieldType": "decimal", "scale": 2},
|
||||
{"name": "decimal_float", "type": ["float", "null"], "fieldType": "decimal", "scale": 2},
|
||||
{"name": "decimal_double", "type": ["double", "null"], "fieldType": "decimal", "scale": 2},
|
||||
{"name": "dateint_bytes_ts", "type": ["bytes", "null"], "fieldType": "dateInt", "layout": "2006-01-02 15:04:05", "unit": "s", "epoch": "1970-01-01 00:00:00"},
|
||||
{"name": "bool_bool", "type": ["boolean", "null"]},
|
||||
{"name": "timestamp_bytes_ts", "type": ["bytes", "null"], "fieldType": "timestamp", "layout": "2006-01-02 15:04:05", "epoch": "1970-01-01 00:00:00"},
|
||||
{"name": "timestamp_bytes_int", "type": ["bytes", "null"], "fieldType": "timestamp", "unit": "s", "layout": "2006-01-02 15:04:05", "epoch": "1970-01-01 00:00:00"}
|
||||
]
|
||||
}
|
||||
14
cli/kafka/testdata/runner/schema/schema02.json
vendored
Normal file
14
cli/kafka/testdata/runner/schema/schema02.json
vendored
Normal file
|
|
@ -0,0 +1,14 @@
|
|||
{
|
||||
"namespace": "org.test",
|
||||
"type": "record",
|
||||
"name": "id_keys",
|
||||
"doc": "All supported avro types and property variations",
|
||||
"fields": [
|
||||
{"name": "pk", "type": "int", "mutex": true, "fieldType": "id"},
|
||||
{"name": "string_string", "type": ["string", "null"], "mutex": true },
|
||||
{"name": "stringset_stringarray", "type": [{"type": "array", "items": "string"}, "null"]},
|
||||
{"name": "idset_longarray", "type": [{"type": "array", "items": "long"}, "null"], "fieldType": "id"},
|
||||
{"name": "int_int", "type": ["int", "null"], "fieldType": "int"},
|
||||
{"name": "bool_bool", "type": ["boolean", "null"]}
|
||||
]
|
||||
}
|
||||
28
idk/Dockerfile-cli-test
Normal file
28
idk/Dockerfile-cli-test
Normal file
|
|
@ -0,0 +1,28 @@
|
|||
ARG GO_VERSION=1.19
|
||||
|
||||
FROM golang:${GO_VERSION} as build_base
|
||||
|
||||
WORKDIR /
|
||||
RUN ["apt-get","update","-y"]
|
||||
RUN ["apt-get","install","-y","git","unixodbc","unixodbc-dev","netcat", "build-essential","musl-tools"]
|
||||
RUN ["git", "clone", "https://github.com/edenhill/librdkafka.git"]
|
||||
WORKDIR /librdkafka
|
||||
RUN ["./configure", "--install-deps"]
|
||||
RUN ["./configure", "--prefix", "/usr"]
|
||||
RUN ["make"]
|
||||
RUN ["make", "install"]
|
||||
|
||||
FROM build_base
|
||||
|
||||
ARG KAFKA_RUNNER_TEST_FEATUREBASE_HOST=pilosa:10101
|
||||
ARG KAFKA_RUNNER_TEST_FEATUREBASEGRPC_HOST=pilosa:20101
|
||||
ARG KAFKA_RUNNER_TEST_KAFKA_HOST=kafka:9092
|
||||
ARG KAFKA_RUNNER_TEST_REGISTRY_HOST=schema-registry:8081
|
||||
|
||||
WORKDIR /go/src/github.com/featurebasedb/featurebase/
|
||||
|
||||
COPY . .
|
||||
|
||||
WORKDIR /go/src/github.com/featurebasedb/featurebase/cli/
|
||||
|
||||
CMD ["go","test","-v","-mod=vendor","-tags=odbc,dynamic" "-run TestKafkaRunner","./..."]
|
||||
16
idk/Makefile
16
idk/Makefile
|
|
@ -338,3 +338,19 @@ update-mocks-%: install-mock-generator
|
|||
$(eval AWS_SDK_VERSION := $(shell grep 'github.com/aws/aws-sdk-go' ../go.mod | cut -d ' ' -f 2))
|
||||
echo Generating mock for AWS service $(AWS_SERVICE) and SDK version $(AWS_SDK_VERSION) && \
|
||||
$(GOPATH)/bin/mockery --name $*API --output idktest/mocks --filename $(AWS_SERVICE).go --dir $(GOPATH)/pkg/mod/github.com/aws/aws-sdk-go@$(AWS_SDK_VERSION)/service/$(AWS_SERVICE)/$(AWS_SERVICE)iface
|
||||
|
||||
|
||||
# Integration Testing for CLI / fbsql tooling that integrates with IDK tooling.
|
||||
test-cli:
|
||||
$(MAKE) startup
|
||||
$(MAKE) test-cli-integration
|
||||
|
||||
TPKG ?= ../cli_test
|
||||
test-cli-integration: testenv vendor
|
||||
$(DOCKER_COMPOSE) build cli-test
|
||||
$(DOCKER_COMPOSE) run -e KAFKA_RUNNER_TEST_FEATUREBASE_HOST=pilosa:10101 \
|
||||
-e KAFKA_RUNNER_TEST_FEATUREBASEGRPC_HOST=pilosa:20101 \
|
||||
-e KAFKA_RUNNER_TEST_KAFKA_HOST=kafka:9092 \
|
||||
-e KAFKA_RUNNER_TEST_REGISTRY_HOST=schema-registry:8081 \
|
||||
-T cli-test bash -c "set -o pipefail; go test -v -mod=vendor -tags=odbc,dynamic -run TestKafkaRunner -covermode=atomic -coverpkg=$(TPKG) -coverprofile=/testdata/$(PROJECT)_base_coverage.out"
|
||||
|
||||
|
|
|
|||
|
|
@ -156,3 +156,19 @@ services:
|
|||
depends_on:
|
||||
- postgres
|
||||
|
||||
cli-test:
|
||||
build:
|
||||
context: ../.
|
||||
dockerfile: ./idk/Dockerfile-cli-test
|
||||
environment:
|
||||
KAFKA_RUNNER_TEST_FEATUREBASE_HOST: pilosa:10101
|
||||
KAFKA_RUNNER_TEST_FEATUREBASEGRPC_HOST: pilosa:20101
|
||||
KAFKA_RUNNER_TEST_KAFKA_HOST: kafka:9092
|
||||
KAFKA_RUNNER_TEST_REGISTRY_HOST: schema-registry:8081
|
||||
volumes:
|
||||
- ./testenv/certs:/certs
|
||||
- ./docker-sasl/ssl_keys:/ssl_keys
|
||||
- ./testdata:/testdata
|
||||
depends_on:
|
||||
- schema-registry
|
||||
|
||||
|
|
|
|||
|
|
@ -659,7 +659,6 @@ initialFetch:
|
|||
for n := range lookupRow {
|
||||
lookupRow[n] = nil
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -690,6 +689,13 @@ initialFetch:
|
|||
|
||||
func (m *Main) Setup() (onFinishRun func(), err error) {
|
||||
if m.basic {
|
||||
// set up SchemaManager which is required for some logic in fbsql
|
||||
// ingest. basic setup doesn't currently support tls so the tls config
|
||||
// is ignored below
|
||||
if _, err = m.setupClient(); err != nil {
|
||||
return nil, errors.Wrap(err, "setting up client")
|
||||
}
|
||||
|
||||
return m.basicSetup()
|
||||
}
|
||||
return m.setup()
|
||||
|
|
|
|||
|
|
@ -23,10 +23,9 @@ func NewMain() (*Main, error) {
|
|||
ConfluentCommand: idk.ConfluentCommand{
|
||||
KafkaBootstrapServers: []string{"localhost:9092"},
|
||||
},
|
||||
Group: "defaultgroup",
|
||||
Topics: []string{"defaulttopic"},
|
||||
Timeout: time.Second,
|
||||
ConsumerCloseTimeout: 30,
|
||||
Group: "defaultgroup",
|
||||
Topics: []string{"defaulttopic"},
|
||||
Timeout: time.Second,
|
||||
}
|
||||
|
||||
m.SchemaRegistryURL = "http://" + defaultRegistryHost
|
||||
|
|
|
|||
|
|
@ -78,12 +78,13 @@ func NewSource() *Source {
|
|||
Group: "group0",
|
||||
Log: logger.NopLogger,
|
||||
|
||||
lastSchemaID: -1,
|
||||
cache: make(map[int32]avro.Schema),
|
||||
recordChannel: make(chan recordWithError),
|
||||
quit: make(chan struct{}),
|
||||
ConfigMap: &confluent.ConfigMap{},
|
||||
highmarks: make([]confluent.TopicPartition, 0),
|
||||
lastSchemaID: -1,
|
||||
cache: make(map[int32]avro.Schema),
|
||||
recordChannel: make(chan recordWithError),
|
||||
quit: make(chan struct{}),
|
||||
ConfigMap: &confluent.ConfigMap{},
|
||||
highmarks: make([]confluent.TopicPartition, 0),
|
||||
consumerCloseTimeout: 30,
|
||||
}
|
||||
|
||||
src.SchemaRegistryURL = "http://" + defaultRegistryHost
|
||||
|
|
@ -224,10 +225,10 @@ func (r *Record) Commit(ctx context.Context) error {
|
|||
p := int32(-1)
|
||||
s := ""
|
||||
r.src.highmarks = r.src.highmarks[:0]
|
||||
|
||||
// sort by increasing partition, decreasing offset
|
||||
|
||||
for _, x := range section {
|
||||
|
||||
if s != *x.Topic || p != x.Partition {
|
||||
r.src.highmarks = append(r.src.highmarks, x)
|
||||
}
|
||||
|
|
@ -235,7 +236,6 @@ func (r *Record) Commit(ctx context.Context) error {
|
|||
s = *x.Topic
|
||||
|
||||
}
|
||||
|
||||
committedOffsets, err := r.src.CommitMessages(r.src.highmarks)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failed to commit messages")
|
||||
|
|
@ -270,7 +270,7 @@ func (s *Source) CommitMessages(recs []confluent.TopicPartition) ([]confluent.To
|
|||
func (s *Source) Open() error {
|
||||
cfg, err := common.SetupConfluent(&s.ConfluentCommand)
|
||||
if err != nil {
|
||||
return err
|
||||
return errors.Wrap(err, "setting up confluent command")
|
||||
}
|
||||
s.ConfigMap = cfg
|
||||
|
||||
|
|
|
|||
|
|
@ -26,9 +26,17 @@ type Source struct {
|
|||
Log logger.Logger
|
||||
Timeout time.Duration
|
||||
SkipOld bool
|
||||
Header string
|
||||
Verbose bool
|
||||
AllowMissingFields bool
|
||||
|
||||
// Header is a file 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
|
||||
|
||||
schema []idk.Field
|
||||
paths idk.PathTable
|
||||
|
||||
|
|
@ -51,11 +59,14 @@ type Source struct {
|
|||
// NewSource gets a new Source
|
||||
func NewSource() *Source {
|
||||
src := &Source{
|
||||
Topics: []string{"test"},
|
||||
Group: "group0",
|
||||
Log: logger.NopLogger,
|
||||
recordChannel: make(chan recordWithError),
|
||||
quit: make(chan struct{}),
|
||||
ConfluentCommand: idk.ConfluentCommand{},
|
||||
Topics: []string{"test"},
|
||||
Group: "group0",
|
||||
Log: logger.NopLogger,
|
||||
recordChannel: make(chan recordWithError),
|
||||
quit: make(chan struct{}),
|
||||
ConfigMap: &confluent.ConfigMap{},
|
||||
highmarks: make([]confluent.TopicPartition, 0),
|
||||
}
|
||||
src.KafkaBootstrapServers = []string{"localhost:9092"}
|
||||
|
||||
|
|
@ -160,21 +171,26 @@ func (r *Record) Commit(ctx context.Context) error {
|
|||
return errors.New("cannot commit a record that has already been committed")
|
||||
}
|
||||
section, remaining := r.src.spool[:idx-base], r.src.spool[idx-base:]
|
||||
// sort by increasing partition, decreasing offset
|
||||
sort.Slice(section, func(i, j int) bool {
|
||||
if section[i].Partition != section[j].Partition {
|
||||
if *section[i].Topic != *section[j].Topic {
|
||||
return *section[i].Topic < *section[j].Topic
|
||||
} else if section[i].Partition != section[j].Partition {
|
||||
return section[i].Partition < section[j].Partition
|
||||
}
|
||||
return section[i].Offset > section[j].Offset
|
||||
})
|
||||
// calculate the high marks
|
||||
p := int32(-1)
|
||||
s := ""
|
||||
r.src.highmarks = r.src.highmarks[:0]
|
||||
|
||||
// sort by increasing partition, decreasing offset
|
||||
for _, x := range section {
|
||||
if p != x.Partition {
|
||||
if s != *x.Topic || p != x.Partition {
|
||||
r.src.highmarks = append(r.src.highmarks, x)
|
||||
}
|
||||
p = x.Partition
|
||||
s = *x.Topic
|
||||
}
|
||||
_, err := r.src.CommitMessages(r.src.highmarks)
|
||||
if err != nil {
|
||||
|
|
@ -193,14 +209,24 @@ 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 (from file or from existing fields)")
|
||||
}
|
||||
|
||||
headerData, err := os.ReadFile(s.Header)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading header file")
|
||||
var headerData []byte
|
||||
var err error
|
||||
if s.Header != "" {
|
||||
headerData, err = os.ReadFile(s.Header)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading header file")
|
||||
}
|
||||
} else {
|
||||
headerData, err = json.Marshal(s.HeaderFields)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling header fields")
|
||||
}
|
||||
}
|
||||
|
||||
schema, paths, err := idk.ParseHeader(headerData)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "processing header")
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue