mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 17:15:56 +00:00
adding cli kafka functionality
This commit is contained in:
parent
aefbc6493b
commit
99c068c261
14 changed files with 343 additions and 75 deletions
|
|
@ -192,7 +192,7 @@ 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:
|
||||
|
|
@ -209,7 +209,7 @@ 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="linux" GOARCH="arm64" make docker-build-fbsql
|
||||
- GOOS="darwin" GOARCH="arm64" make docker-build-fbsql
|
||||
artifacts:
|
||||
paths:
|
||||
|
|
@ -342,6 +342,22 @@ run go tests dax/test/dax:
|
|||
paths:
|
||||
- coverage-dax-integration.out
|
||||
|
||||
# fbsql test
|
||||
run go tests cli kafka integration:
|
||||
stage: integration
|
||||
image: golang:$GOVERSION
|
||||
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
|
||||
script:
|
||||
- echo "running fbsql integration tests"
|
||||
- cd ./idk/
|
||||
- BRANCH_NAME=${CI_COMMIT_REF_SLUG} make start-all
|
||||
- cd ..
|
||||
- go test -run TestKafkaRunner ./cli
|
||||
|
||||
# idk tests
|
||||
run go tests idk race:
|
||||
variables:
|
||||
|
|
|
|||
58
Dockerfile-darwin-cgo-builder
Normal file
58
Dockerfile-darwin-cgo-builder
Normal file
|
|
@ -0,0 +1,58 @@
|
|||
# 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.
|
||||
#
|
||||
# For this to work, that Xcode.xip file must be in the working directory. For
|
||||
# now, you'll need to manually define the version.
|
||||
|
||||
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/*
|
||||
|
||||
# build compilers for cross compilation (darwin on linux) with cgo
|
||||
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 \
|
||||
79
Dockerfile-fbsql-darwin
Normal file
79
Dockerfile-fbsql-darwin
Normal file
|
|
@ -0,0 +1,79 @@
|
|||
# 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 \
|
||||
&& 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
|
||||
|
|
@ -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
|
||||
46
Makefile
46
Makefile
|
|
@ -361,38 +361,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 -tags dynamic $(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)
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ import (
|
|||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
featurebase "github.com/featurebasedb/featurebase/v3"
|
||||
"github.com/featurebasedb/featurebase/v3/cli"
|
||||
|
|
@ -127,19 +128,42 @@ var kafkaRunnerTests = []kafkaRunnerTest{
|
|||
},
|
||||
CreateTableStmt: "(_id String, name String, age Int, hobbies StringSet)",
|
||||
},
|
||||
/*{ // id keys
|
||||
{ // string keys
|
||||
ConfigFile: "config04.toml",
|
||||
DataFile: "data01.json",
|
||||
Encode: "avro",
|
||||
SchemaFile: "schema01.json",
|
||||
Tests: []testQuery{
|
||||
{
|
||||
Query: "Extract(Sort(All(), field=name), Rows(name), Rows(age), Rows(hobbies))",
|
||||
ExpectedResp: `{"results":[{"fields":[{"name":"name","type":"string"},{"name":"age","type":"int64"},{"name":"hobbies","type":"[]string"}],"columns":[{"column":1,"rows":["a",20,["hob2","hob1"]]},{"column":2,"rows":["b",21,["hob2","hob3"]]},{"column":3,"rows":["c",22,["hob3","hob4"]]},{"column":4,"rows":["d",23,["hob4","hob5"]]},{"column":5,"rows":["e",24,["hob5","hob6"]]},{"column":6,"rows":["f",26,["hob6","hob7"]]}]}]}`,
|
||||
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
|
||||
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
|
||||
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.
|
||||
|
|
@ -151,9 +175,9 @@ var kafkaRunnerTests = []kafkaRunnerTest{
|
|||
// 5. Runs the cli.Command
|
||||
// 6. Confirms that the data was written to FeatureBase as expected
|
||||
func TestKafkaRunner(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test")
|
||||
}
|
||||
//if testing.Short() {
|
||||
// t.Skip("skipping integration test")
|
||||
//}
|
||||
|
||||
// Get host and port pair for services required for test (e.g. kafka and
|
||||
// featurebase)
|
||||
|
|
@ -189,6 +213,7 @@ func TestKafkaRunner(t *testing.T) {
|
|||
|
||||
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
|
||||
|
||||
|
|
@ -234,7 +259,7 @@ func TestKafkaRunner(t *testing.T) {
|
|||
schema := string(schemaBytes)
|
||||
|
||||
// post the schema to schema registry
|
||||
schemaID, err := postSchema(schema, "kafka-runner-subject", services.registryHost)
|
||||
schemaID, err := postSchema(schema, "kafka-runner-subject-"+time.Now().String(), services.registryHost)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -309,26 +334,26 @@ func createTempFindAndReplace(source string, mapping map[string]string) error {
|
|||
// 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 [][]interface{}
|
||||
for _, line := range wqr.Data {
|
||||
var newline []interface{}
|
||||
for _, element := range line {
|
||||
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:
|
||||
newline = append(newline, newElement.SortedStringSlice())
|
||||
newline[j] = newElement.SortedStringSlice()
|
||||
case featurebase.IDSet:
|
||||
newline = append(newline, newElement.SortedInt64Slice())
|
||||
newline[j] = newElement.SortedInt64Slice()
|
||||
default:
|
||||
newline = append(newline, element)
|
||||
newline[j] = element
|
||||
}
|
||||
}
|
||||
data = append(data, newline)
|
||||
data[i] = newline
|
||||
}
|
||||
js, err := json.Marshal(data)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
//t.Fatal(string(js))
|
||||
// t.Fatal(string(js))
|
||||
require.JSONEq(t, expectedQuery, string(js))
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -49,6 +49,11 @@ func (cmd *Command) newKafkaRunner(cfgFile string) (*kafka.Runner, error) {
|
|||
return nil, errors.Wrap(err, "getting fields from config")
|
||||
}
|
||||
|
||||
// for avro, let the SchemaManager and IDK handle fields
|
||||
if cfg.Encode == "avro" {
|
||||
flds = nil
|
||||
}
|
||||
|
||||
return kafka.NewRunner(
|
||||
idkCfg,
|
||||
batch.NewSQLBatcher(cmd, flds),
|
||||
|
|
|
|||
|
|
@ -41,7 +41,6 @@ func NewRunner(cfg ConfigForIDK, batcher fbbatch.Batcher, logWriter io.Writer) *
|
|||
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,
|
||||
Hosts: cfg.Hosts,
|
||||
|
|
|
|||
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"
|
||||
schemaRegistryHost = "localhost:8081"
|
||||
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"
|
||||
schemaRegistryHost = "localhost:8081"
|
||||
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
|
||||
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}}
|
||||
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"]}
|
||||
]
|
||||
}
|
||||
|
|
@ -301,42 +301,59 @@ func (m *Main) run() error {
|
|||
|
||||
func (m *Main) clone() (*Main, error) {
|
||||
var index *pilosaclient.Index
|
||||
noOpSchemaManager := false
|
||||
|
||||
// If you have a schema manager, it does it's thing. Otherwise, it's a no
|
||||
// opt manager. Then you get a default schema. If you get a default schema,
|
||||
// you get a default index. This seems fine for IDK. However, for the CLI /
|
||||
// SQL kafka runner, we use the m.index to create the dax.Table. If m.index
|
||||
// is set to default values, then keys is false even when we don't want it
|
||||
// to be. This was causing issue in buildBulkInsert in the batch package of
|
||||
// the CLI.
|
||||
switch m.SchemaManager.(type) {
|
||||
case *nopSchemaManager:
|
||||
noOpSchemaManager = true
|
||||
}
|
||||
|
||||
schema, err := m.SchemaManager.Schema()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if noOpSchemaManager && len(m.PrimaryKeyFields) > 0 {
|
||||
// if IDK has PrimaryKeyFields, it expects keyed index
|
||||
keys := pilosaclient.OptIndexKeys(true)
|
||||
// most queries don't work with this set to false so set to true
|
||||
exists := pilosaclient.OptIndexTrackExistence(true)
|
||||
index = schema.Index(m.Index, keys, exists)
|
||||
} else {
|
||||
index = schema.Index(m.Index)
|
||||
}
|
||||
index = schema.Index(m.Index)
|
||||
|
||||
// use a copy (schema race condition issues)
|
||||
// use a copy (schema race condition issu)
|
||||
mClone := *m
|
||||
mClone.index = index
|
||||
|
||||
return &mClone, nil
|
||||
}
|
||||
|
||||
// func (m *Main) clone() (*Main, error) {
|
||||
// var index *pilosaclient.Index
|
||||
// noOpSchemaManager := false
|
||||
|
||||
// // If you have a schema manager, it does it's thing. Otherwise, it's a no
|
||||
// // opt manager. Then you get a default schema. If you get a default schema,
|
||||
// // you get a default index. This seems fine for IDK. However, for the CLI /
|
||||
// // SQL kafka runner, we use the m.index to create the dax.Table. If m.index
|
||||
// // is set to default values, then keys is false even when we don't want it
|
||||
// // to be. This was causing issue in buildBulkInsert in the batch package of
|
||||
// // the CLI.
|
||||
// switch m.SchemaManager.(type) {
|
||||
// case *nopSchemaManager:
|
||||
// noOpSchemaManager = true
|
||||
// }
|
||||
|
||||
// schema, err := m.SchemaManager.Schema()
|
||||
// if err != nil {
|
||||
// return nil, err
|
||||
// }
|
||||
|
||||
// if noOpSchemaManager && len(m.PrimaryKeyFields) > 0 {
|
||||
// // if IDK has PrimaryKeyFields, it expects keyed index
|
||||
// keys := pilosaclient.OptIndexKeys(true)
|
||||
// // most queries don't work with this set to false so set to true
|
||||
// exists := pilosaclient.OptIndexTrackExistence(true)
|
||||
// index = schema.Index(m.Index, keys, exists)
|
||||
// } else {
|
||||
// index = schema.Index(m.Index)
|
||||
// }
|
||||
|
||||
// // use a copy (schema race condition issues)
|
||||
// mClone := *m
|
||||
// mClone.index = index
|
||||
|
||||
// return &mClone, nil
|
||||
// }
|
||||
|
||||
func (m *Main) runIngester(c int, l *msgCounter) error {
|
||||
m.log.Printf("start ingester %d", c)
|
||||
// TODO: actually implement cancellation and graceful shutdown
|
||||
|
|
@ -727,6 +744,8 @@ func (m *Main) basicSetup() (onFinishRun func(), err error) {
|
|||
return nil, errors.Wrap(err, "validating configuration")
|
||||
}
|
||||
|
||||
_, err = m.setupClient()
|
||||
|
||||
// setup logging
|
||||
var f *logger.FileWriter
|
||||
var logOut io.Writer = os.Stderr
|
||||
|
|
@ -2184,7 +2203,7 @@ func (m *Main) newBatch(clientFields []*pilosaclient.Field) (pilosabatch.RecordB
|
|||
ii := pilosaclient.FromClientIndex(m.index)
|
||||
tbl := pilosacore.IndexInfoToTable(ii)
|
||||
|
||||
// Fields.
|
||||
// Fields. // this is giving bad fieldInfos
|
||||
fields := pilosaclient.FromClientFields(clientFields)
|
||||
|
||||
// If a custom Batcher has been defined, use that. Otherwise default to
|
||||
|
|
|
|||
|
|
@ -215,9 +215,7 @@ func (ii IDSet) String() string {
|
|||
// sorted
|
||||
func (ii IDSet) SortedInt64Slice() []int64 {
|
||||
var idSetSlice = make([]int64, len(ii))
|
||||
for i, str := range ii {
|
||||
idSetSlice[i] = str
|
||||
}
|
||||
copy(idSetSlice, ii)
|
||||
sort.Slice(idSetSlice, func(i, j int) bool { return idSetSlice[i] < idSetSlice[j] })
|
||||
return idSetSlice
|
||||
}
|
||||
|
|
@ -243,9 +241,7 @@ func (ss StringSet) String() string {
|
|||
// that is sorted
|
||||
func (ss StringSet) SortedStringSlice() []string {
|
||||
var stringSetSlice = make([]string, len(ss))
|
||||
for i, str := range ss {
|
||||
stringSetSlice[i] = str
|
||||
}
|
||||
copy(stringSetSlice, ss)
|
||||
sort.Strings(stringSetSlice)
|
||||
return stringSetSlice
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue