Kafka Streams and AI
First written October 2025, last updated September 2026.
This walkthrough covers running large language models (LLMs) on your laptop for local development and using them with Kafka Streams for state store analysis.
Goals
- Set up a local LLM for application development
- Integrate the LLM with Kafka Streams for state store analysis
When using a local LLM, your performance will depend on your hardware. We used an M1 Max MacBook Pro.
Setting Up a Local LLM
Our first attempt was deploying Ollama inside a Docker container. While we expected some slowness, we hoped to tune Docker to access the GPU, but this wasn’t straightforward. So we shifted to a local installation: we installed Ollama via Homebrew, downloaded the model, and used two hostnames to access it:
localhostfor development on the laptophost.docker.internalfor accessing localhost services from within Docker containers
Installation
All steps below assume macOS with Homebrew installed.
Install Ollama
Installation is straightforward with Homebrew.
brew install ollamaDownload a Model
The hardest part of downloading a model is selecting one. Once selected, you download it like a Docker image:
ollama pull deepseek-r1:latestStart Ollama Server
You can run Ollama manually or as a service.
Manually:
ollama serveAs a service:
brew services start ollamabrew services stop ollamaWe recommend manual execution, as the console log is very insightful. It's easy to change settings by setting environment variables. Here’s a simple way to increase context length:
OLLAMA_CONTEXT_LENGTH=16384 ollama serveModel Selection
Browse available models at ollama.com/library. We used deepseek-r1 for its reasoning capability and laptop-friendly sizes; pick whichever current model fits your hardware. Once local development works, migrating to a cloud environment becomes simpler. After pulling multiple models, use ollama list to view them.
Endpoint Selection
Ollama endpoint
While we started with the standard Ollama chat endpoint (/api/chat), we also experimented with the OpenAI endpoint (/v1/chat/completions). That compatibility layer makes it easy to point existing OpenAI-client code at a local model. On the native endpoint, sampling settings such as temperature and top_k go in the request's options object.
Endpoint configuration
Java records and Jackson for serialization are an easy and powerful way to quickly work with a REST API from your Java application.
record Message(String role, String content) {}record OllamaRequest( String model, Boolean stream, List<Message> messages, Map<String, Object> options) {}record OllamaResponse( String model, String created_at, Message message, Boolean done, String done_reason, long total_duration, long load_duration, long prompt_eval_count, long prompt_eval_duration, long eval_count, long eval_duration) {}HTTP Client
Use whatever client you are currently using from your Java applications; we’ve been quite happy using Apache HttpComponents, HttpClient.
Integrating the LLM with Kafka Streams
With Ollama running and a model downloaded, it was time to integrate with Kafka Streams.
Kafka Streams excels at stateful stream processing, but it isn't designed for non-stream computation inside stream threads; doing so hurts performance.
Kafka Streams exposes state stores via interactive queries:
streams.store()(IQv1)streams.query()(IQv2)
If your state store is analytical in nature, it is a great data source for AI. Do not execute queries to the model from within the stream thread (topology).
The code is simple. The prompts took the work: learning how to write, test, and refine them to improve the analysis.
The demo application and LLM integration are in the kafka-streams-dashboards project.
Accessing the State Stores
To provide context in prompts to the LLM, include attributes for the products that make the top 5. Global stores are handy for this.
Use the IQv1 API to access any global store. There is currently no IQv2 API for global stores. For this walkthrough, global stores are used for product information. This can run on a schedule or via user interaction. For demonstration, it is scheduled to run at the end of each window.
var products = new HashMap<>();var q = streams.store(StoreQueryParameters.fromNameAndType("product", QueryableStoreTypes.keyValueStore()));try (var iterator = q.all()) { while (iterator.hasNext()) { var next = iterator.next(); products.put(next.key, next.value); }}The data we want the LLM to analyze resides in tumbling window state stores.
Use the IQv2 API to access a specific window from the state store. While the IQv1 API is still available, use the IQv2 API when possible. For a tumbling window, use a range query with the start/end aligned to the tumbling window definition.
Create a top-N list for easy consumption when writing the prompt.
List<ProductAnalytic> pa = new ArrayList<>(); Window current = currentWindow(windowSize); var request = StateQueryRequest.inStore("TUMBLING-aggregate-purchase-order") .withQuery(WindowRangeQuery.withWindowStartRange(current.start(), current.end()));var result = streams.query(request); for (var partitionResult : result.getPartitionResults().values()) { if (partitionResult.isSuccess()) { try (var iterator = partitionResult.getResult()) { while (iterator.hasNext()) { var next = iterator.next(); Windowed<String> windowedKey = next.key; ProductAnalytic value = next.value.value(); // the store returns TimestampAndValue; value.value() unwraps it pa.add(value); } } }} pa.sort(Comparator.comparingLong(ProductAnalytic::getQuantity).reversed());Writing a prompt is easy. Writing one that gets consistent answers took most of our time. Once orchestration is in place, collect your top-N datasets and construct the prompt. Expect weak outputs at first, and improve them as you learn how to communicate with your model effectively. If you want to build on context from past responses, confirm that process is part of your model testing and analysis.
var sb = new StringBuilder();sb.append("predict trends on purchasing, considering previous window purchases\n");sb.append("top-5 quantity for current window:\n");sb.append(String.format("window [%s, %s)\n", current.start(), current.end()));pa.subList(0, Math.min(pa.size(), 5)).forEach(p -> { var product = products.get(p.getSku()); sb.append(String.format("sku=%s, qty=%d, attrs=%s\n", p.getSku(), p.getQuantity(), product.attributes()));});observer.inquire(sb.toString());The Observer retains past responses and system messages, and drops the user turns that produced them, to keep input tokens lower. The tradeoff: the model sees its own earlier answers without the data they were about, which may explain some of the inconsistency below. Sending both turns, or a summary of each exchange, is the first thing to try instead.
public Observer() { httpClient = HttpClients.custom() .setDefaultRequestConfig(RequestConfig.custom() .setResponseTimeout(Timeout.ofSeconds(60)) .build() ) .build(); messages.add(new Message(SYSTEM, "the time windows are small for demo purposes, so don't worry about that; it is correct.")); messages.add(new Message(SYSTEM, "focus on trends, predict future purchases")); messages.add(new Message(SYSTEM, "You are a concise assistant. Do not include any <think> sections or internal reasoning in your responses. Only return direct answers or summaries.")); }public void inquire(String prompt) throws IOException { List<Message> request = new ArrayList<>(this.messages); request.add(new Message("user", prompt)); HttpPost post = new HttpPost(ollamaUrl); post.setHeader("Accept", "application/json"); post.setHeader("Content-Type", "application/json"); post.setEntity(new StringEntity( objectMapper.writeValueAsString( new OllamaRequest( "deepseek-r1", false, request, null ) ), ContentType.APPLICATION_JSON) ); OllamaResponse response = httpClient.execute(post, (ClassicHttpResponse resp) -> { final int status = resp.getCode(); if (status >= 200 && status < 300) { return objectMapper.readValue(resp.getEntity().getContent(), OllamaResponse.class); } throw new RuntimeException("status=" + status); }); //store response to be provided for the next inquiry, trim/filter as needed. messages.add(response.message()); }Architecture highlights
The architecture highlights are:
- The prompt is a short write-up with the time window and top-5 products, including all product attributes.
- The model is
deepseek-r1running in Ollama. - Kafka Streams uses IQv1 and IQv2 APIs to read state stores from outside the topology.
- The LLM receives current information and past responses.
- Only historical responses were provided in the next conversation, not the historical requests.
Walkthrough
Here is a walkthrough of our example. It shows the importance of making your LLM conversations replayable, so you can tune and verify.
Actual Interaction
Orders are created once every second, with a publisher that randomly selects products from the product list.
./publisher/run.shAfter some analysis, restart the publisher with a bias toward selecting tablets more frequently than phones and computers.
./publisher/run.sh --attribute-name productType --attribute-value tablets --percentage 0.80Response 1
- Work is needed to reduce thinking content in responses. The durable fix is Ollama's think option: set think to false in the request and reasoning output is suppressed at the source. Strip whatever still leaks through.
- Little perceived value in this response.
Message[role=assistant, content=<think>Okay, the user is asking to predict trends in purchasing based on previous window data. They provided five entries with SKUs and their quantities within a time window from 2025-10-22T11:55:30Z to 2025-10-22T11:56:00Z, which is very small.[reasoning trimmed]</think>Based on these very recent purchases (a micro-window), trends observed are: 1. Higher quantities purchased (`qty=5` or `6`) for items with larger RAM size (`ramSize=12GB`, `16GB`) and higher review counts/scores (`reviewCount >=900`, `reviewScore>=4.1`).2. Items from newer product lines (like tablets) might be seeing slightly more purchase volume in this period.3. The specific CPU model (`M2`) continues to influence purchases within computers and tablets seen here.Response 2
- While data is random, product attributes show patterns, and the model appears to detect similarities.
- The response provided a prediction of future trends, indicating our system request “focus on trends, predict future purchases” was understood.
Based on these recent purchases (a very short time window), the observed trend is: * **Higher quantities purchased:** `qty` ranges from 2 to 7, with most entries being higher than average.* **Dominant product type:** Most items are computers (`productType=computers`).[...] **Prediction for the next similar small time window:** * Computers with large amounts of RAM (>= 16GB) and fast SSDs will likely continue to be popular.[...]Response 3
- The model didn't provide a prediction, and the volumes it reports (hundreds of thousands of purchases in a 10-second window from a once-per-second producer) are invented, along with the product names. Concrete-sounding numbers still need validation against the store.
In this short time window (10 seconds), purchase activity remains dominated by SKUs from the previous top-5 list. * **SKU 103** ("Apple MacBook Pro") continues to be a massive driver of volume (~248k+).* **SKU 105** ("Samsung Galaxy S Series Tablet - M2 Chip") leads this period with over 297k purchases.* **SKU 129** ("Dell XPS Laptop - New Model") debuts strongly, securing the third position (~50k+).Producer Changed
- The producer was redeployed with an 80% bias toward tablets.
Response 4
- A prediction is provided; the change in product distribution isn't reflected.
- The formatting is inconsistent; this is why we’re considering fine-tuning a model instead to see if we can get more consistent responses.
**Observed Trends in Current Window:** 1. **Dominance of Electronics:** Top purchases are exclusively from `tablets` and `computers`, alongside `phones`.2. **Lower Purchase Volume for Computers:** The highest quantity (`qty`) recorded is `4`, significantly lower than the previous window's peak (which was higher, e.g., around `50k+`). Only one computer SKU appears in this top-5.[...] **Predicted Trend for Future:** * **Continued Focus on Tablets/Phones:** These categories are expected to maintain high purchase volumes in future short time windows due to their presence and performance here.[...]Response 5
- Very short response.
- Again, response formats are wildly different, even though the inputs are uniform.
Based on this current top-5 (all tablets) with similar attributes (M2/A16/A17 chips, high storage/RAM), the trend remains focused on **tablets** from specific brands like **Apple** and one **Samsung**,driven by strong performance specifications. Purchase quantities are lower than previous windows but still reflect sustained interest in these products during short time periods.Response 6
- The LLM provides the insight we hoped it would see: “High demand is expected to continue for tablets.”
- It also suggests “high review counts/scores will continue to be popular.” This could be due to the sample skew or inference from limited data.
Based on the top-5 purchased SKUs from the current very short time window: 1. **Dominant Product Type:** Tablets significantly dominate this purchase period (4 out of 5 items).2. **Key Attributes:** Items feature high SSD storage, LPDDR5 RAM, and good review scores/counts.3. **Performance Specs:** High CPU specs (e.g., `cpuCores=10`, `gpuCores=12`) are associated with higher quantities (`qty=2`). Lower quantities (`qty=1` or `3`) correspond to items with slightly lower RAM (`ramSize=8GB`), older release dates, or different screen ratios.4. **Stable Tablet Focus:** This reinforces the trend observed previously that tablets remain a primary focus for high purchase volume in short time windows. **Predicted Trend:** Based *strictly* on these recent purchases: * High demand is expected to continue for **tablets**, particularly those with **high performance (e.g., many CPU/GPU cores) and large amounts of RAM/SSD storage.*** Newer tablet models (`released`=2025-01-28, 2025-04-13) are likely to maintain strong purchase interest.* Items with **high review counts/scores** will continue to be popular.Conclusion
The pattern works: Kafka Streams state is a clean input for a model, read through interactive queries from outside the stream threads. The model was the weak link. Six responses gave inconsistent formats, one invented purchase volumes and product names, and the one real insight (tablets) arrived only after the producer was biased toward them. What we'd change now: disable reasoning output at the source, pin a response format in the system prompt, and send the user turns along with the responses so the model sees the whole conversation.
Takeaways
- It's possible to use local LLMs for developer experience, exploration, and functional/integration tests.
- The samples here use Ollama's native
/api/chat. Use the OpenAI-compatible endpoint (/v1/chat/completions) when you want to swap models or providers without changing client code. - Put in the effort to alter and replay conversations when experimenting with the best way to interact with LLMs.
- Disable reasoning output at the source where the API supports it (Ollama's think option), and prune whatever still leaks through.
- Capture your data so you can alter prompts and repeat.
- The better you know the data trends yourself, the easier it is to validate whether the LLM is seeing what you expect, and to notice when it finds something you didn't expect.
Next Steps
Our next considerations are:
- Combining state stores across Kafka Streams instances for LLM analysis.
- Tokenization analysis to reduce token spend.
- Fine-tuning a model with product information and comparing the base model with the tuned model.
Working on something like this?
Start a Conversation