When there are multiple reference-style markdown links in the same deck with the same label, they will silently clash - i.e. one will overwrite the other. The problem can become very apparent when using many links like [see the docs][docs] in different slides, where [docs] points to a different URL each time. This commit adds a crude script to detect such duplicates and display them. This script was used to detect a bunch of duplicates and fix them (by making the label unique). There are still a few duplicates left but they point to the same places, so we decided to leave them as-is for now (but might change that later).
12 KiB
Getting started with Bento
How can we move to a message queue architecture...
...without rewriting a bunch of code?
🤔
Bento
"Fancy stream processing made operationally mundane"
"Written in Go, deployed as a static binary, declarative configuration. Open source and cloud native as utter heck."
With ✨ amazing ✨ documentation 😍
class: extra-details
Tiny bit of history
-
Original project: Benthos
-
May 30, 2024: Redpanda acquires Benthos
-
Benthos is now Redpanda Connect
-
some parts have been relicensed as commercial products
-
-
May 31, 2024: Warpstream forks Benthos
-
that fork is named "Bento"
-
it's fully open source
-
-
We're going to use Bento here, but Redpanda Connect should work fine too!
Bento concepts
-
Message stream processor
-
Each pipeline is configured by a YAML configuration that defines:
-
input (where do we get the messages?)
-
pipeline (optional: how do we transform the messages?)
-
output (where do we put the messages afterwards?)
-
-
Once Bento is started, it runs the pipelines forever
(except for pipelines that have a logical end, e.g. reading from a file)
-
Embedded language (Bloblang) to manipulate/transform messages
Messages
-
Typically JSON objects
(but raw strings are also possible)
-
Nesting, arrays, etc. are OK
Getting started with Bento
We're going to:
-
Import a bunch of cities from a CSV file into a Redis queue.
-
Read back these cities using a web server.
-
Use an "enrichment workflow" to query our LLM for each city.
1️⃣ Importing cities
Let's break down the work:
-
download the data set
-
create the Bento configuration
-
deploy Redis
-
start Bento
Downloading the data set
-
Example database:
-
Let's download and uncompress the data set:
curl -fsSL https://www.kaggle.com/api/v1/datasets/download/juanmah/world-cities | funzip > cities.csv(Ignore the "length error", it's harmless!)
-
Check the structure of the data set:
head cities.csv
Creating the Bento configuration
-
We need to find which
inputandoutputto use -
Check the list with
bento listor the documentation -
Then run
bento create INPUTNAME/PIPELINENAME/OUTPUTNAME -
Generate a configuration file:
bento create csv//redis_list > csv2redis.yaml -
Edit that configuration file; look for the
(required)parameters(Everything else can go away!)
Resulting configuration
If we trim all the default values, here is the result:
input:
csv:
paths: ["cities.csv"]
output:
redis_list:
url: redis://redis:6379 # No default (required)
key: cities
We'll call that value csv2redis.yaml.
Deploying Redis
-
Create a Deployment:
kubectl create deployment redis --image redis -
Expose it:
kubectl expose deployment redis --port 6379
Starting Bento
Option 1: run it manually in a pod, to see what's going on.
bento --config csv2redis.yaml
Option 2: run it with e.g. the Bento Helm chart.
We're not going to do that yet, since this particular pipeline has a logical end.
(The Helm chart is best suited to pipelines that run forever.)
Expected output
.small[
INFO Running main config from specified file @service=bento bento_version="" path=csv2redis.yaml
INFO Launching a Bento instance, use CTRL+C to close @service=bento
INFO Listening for HTTP requests at: http://0.0.0.0:4195 @service=bento
INFO Input type csv is now active @service=bento label="" path=root.input
INFO Output type redis_list is now active @service=bento label="" path=root.output
INFO Pipeline has terminated. Shutting down the service @service=bento
]
The pipeline should complete in just a few seconds.
Checking what's in Redis
-
Connect to our Redis instance:
redis-cli -h redis -
List keys:
KEYS * -
Check that the
citieslist has approx. 47000 elements:LLEN cities -
Get the first element of the list:
LINDEX cities 0
Fun with Bloblang
-
Let's add a filter to keep only cities with a population above 10,000,000
-
Add the following block to the Bento configuration:
pipeline:
processors:
- switch:
- check: this.population == ""
processors:
- mapping: root = deleted()
- check: this.population.int64() < 10000000
processors:
- mapping: root = deleted()
(See the docs for details about the switch processor.)
Testing our processor
-
First, delete the existing
citieslist:redis-cli -h redis DEL cities -
Then, run the Bento pipeline again:
bento --config csv2redis.yaml(It should complain about a few cities where the population has a decimal point.)
-
Check how many cities were loaded:
redis-cli -h redis LLEN cities(There should be 47.)
2️⃣ Consume the queue over HTTP
-
We want to "get the next city" in the queue with a simple
curl -
Our input will be
redis_list -
Our output will be
http_server
Generate the Bento configuration
Option 1: bento create redis_list//http_server
Option 2: read the docs
🙋 Choose your own adventure
Do you want to try to write that configuration?
Or shall we see it right away?
--
⚠️ Spoilers on next slide!
redis2http.yaml
input:
redis_list:
url: redis://redis:`6379`
key: cities
output:
http_server:
path: /nextcity
This will set up an HTTP route to fetch one city.
It's also possible to batch, stream...
⚠️ As of November 2024, bento create uses port 6397 instead of 6379 for Redis!
Trying it out
-
Run Bento with this configuration:
bento --config redis2http.yaml & -
Retrieve one city:
curl http://localhost:4195/nextcity -
Check what happens after we retrive all the cities!
3️⃣ Query our LLM for each city
-
We want to ask our LLM who's the mayor of each of these cities
-
We'll use a prompt that will usually ensure a short answer
(so that it's faster; we don't want to wait 30 seconds per city!)
-
We'll test the prompt with the Ollama CLI
-
Then we'll craft a proper HTTP API query
-
Finally, we'll configure an enrichment workflow in Bento
Test our prompt
Assuming that our earlier Ollama Deployment is still running:
kubectl exec deployment/ollama -- \
ollama run qwen2:1.5b "
Who is the mayor of San Francisco?
Just give the name by itself on a single line.
If you don't know, don't say anything.
"
Turn the prompt into an HTTP API query
Note: to install http in an Alpine container, run apk add httpie.
http http://ollama.default:11434/api/generate \
model=qwen2:1.5b stream:=false prompt="
Who is the mayor of Paris?
Just give the name by itself on a single line.
If you don't know, don't say anything.
"
We get a JSON payload, and we want to use the response field.
Configure an enrichment workflow
The Bento documentation is really good!
We need to set up:
-
a
branchprocessor -
a
request_mapto transform the city into an Ollama request -
an
httpprocessor to submit the request to Ollama -
a
result_mapto transform the Ollama response
Without the branch processor
flowchart LR
CITY["
city: Paris
country: France
population: 1106000
iso2: FR
...
"]
REQ["
model: qwen2:1.5b
stream: false
prompt: Who is the mayor of Paris?
"]
REP["
response: Anne Hidalgo
eval_count: ...
prompt_eval_count: ...
(other ollama fields)
"]
CITY@{ shape: card}
REQ@{ shape: card}
REP@{ shape: card}
style CITY text-align: left
style REQ text-align: left
style REP text-align: left
mapping@{ shape: diam }
http["http processor"]@{ shape: diam }
CITY --> mapping --> REQ --> http --> REP
-
We transform the
cityinto an Ollama request -
The
httpprocessor submits the request to Ollama -
The final output is the Ollama response
With the branch processor
flowchart LR
CITY["
city: Paris
country: France
population: 1106000
iso2: FR
...
"]
REQ["
model: qwen2:1.5b
stream: false
prompt: Who is the mayor of Paris?
"]
REP["
response: Anne Hidalgo
eval_count: ...
prompt_eval_count: ...
(other ollama fields)
"]
OUT["
city: Paris
country: France
population: 1106000
iso2: FR
...
mayor: Anne Hidalgo
"]
CITY@{ shape: card}
REQ@{ shape: card}
REP@{ shape: card}
OUT@{ shape: card}
style CITY text-align: left
style REQ text-align: left
style REP text-align: left
style OUT text-align: left
branch@{ shape: diam }
request_map@{ shape: diam }
result_map@{ shape: diam }
http["http processor"]@{ shape: diam }
CITY --> branch
branch --> result_map
branch --> request_map
request_map --> REQ
REQ --> http
http --> REP
REP --> result_map
result_map --> OUT
-
The
branchprocessor allows doing the processing "on the side" -
request_mapandresult_maptransform the message before/after processing -
Then, the result is combined with the original message (the
city)
input:
csv:
paths: ["cities.csv"]
pipeline:
processors:
- branch:
request_map: |
root.model = "qwen2:1.5b"
root.stream = false
root.prompt = (
"Who is the mayor of %s? ".format(this.city) +
"Just give the name by itself on a single line. " +
"If you don't know, don't say anything."
)
processors:
- http:
url: http://ollama:11434/api/generate
verb: POST
result_map: |
root.mayor = this.response
Trying it out
-
Save the YAML on the previous page into a configuration file
-
Run Bento with that configuration file
-
What happens?
--
🤔 We're seeing errors due to timeouts
ERRO HTTP request to 'http://ollama...' failed: http://ollama...:
Post "http://ollama...": context deadline exceeded
(Client.Timeout exceeded while awaiting headers)
🙋 Choose your own adventure
How should we address errors?
-
Option 1: increase the timeout in the http processor
-
Option 2: use a retry processor in the pipeline
-
Option 3: use a reject_errored output
🏗️ Let's build something!
-
We want to process 1000 cities with our LLM
(guessing who the mayor is, or something similar)
-
Store the output wherever we want
(Redis, CSV file, JSONL files...)
-
Deal correctly with errors
(we'll check that there are, indeed, 1000 cities in the output)
-
Scale out to process faster
(scale ollama to e.g. 10 replicas, enable parallelism in Bento)
class: title
🍱 Lunch time! 🍱
What happened?
-
If your Ollama pods have resource requests:
→ your cluster may have auto-scaled
-
If your Ollama pods don't have resource requests:
→ you probably have a bunch of container restarts, due to out-of-memory errors
🤔 What's that about?