aws-samples / aws-samples/sample-stream-ingest-amazon-bedrock-knowledge-base

Use MSK CreateTopic API, upgrade Kafka version, and replace pykafka with confluent-kafka

Open
#3 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Jupyter Notebook
Stars
0
Forks
1
PR merge metrics
No merged PRs in 30d

Description

## Summary

This issue proposes several improvements to simplify the setup process and replace deprecated dependencies.

## Proposed Changes

### 1. Use MSK CreateTopic API instead of manual Kafka client installation

The current `1.Setup.ipynb` requires users to open a SageMaker Studio terminal and manually install Java, wget, Apache Kafka client, and aws-msk-iam-auth JAR to create a Kafka topic. This is error-prone and time-consuming.

Amazon MSK now provides [public APIs for Kafka topic management](https://aws.amazon.com/about-aws/whats-new/2026/02/amazon-msk-kafka-topics-public-apis/) (CreateTopic, ListTopics, etc.), available for clusters running Kafka 3.6+. Using these APIs from the notebook eliminates the need for terminal commands and Kafka client installation entirely.

**Before:** 7-step terminal procedure (install Java, wget, download Kafka client, download IAM auth JAR, create topic, verify, cleanup)
**After:** 2 cells in the notebook (create_topic + list_topics via boto3)

### 2. Upgrade MSK Kafka version from 2.8.0 to 3.8.x

The CloudFormation template currently uses Kafka 2.8.0. Upgrading to 3.8.x is required for the CreateTopic API (needs Kafka 3.6+) and is the [AWS recommended version](https://docs.aws.amazon.com/msk/latest/developerguide/supported-kafka-versions.html).

### 3. Replace pykafka with confluent-kafka in 2.StreamIngest.ipynb

[PyKafka](https://github.com/Parsely/pykafka) was archived on 2021-03-24 and its last release was in 2018. The [confluent-kafka](https://pypi.org/project/confluent-kafka/) library is actively maintained (v2.13.0 as of January 2026) and provides Kafka 3.9.x compatibility through librdkafka.

### 4. Add CloudFormation Parameters and Outputs for KB/DS names

Add KnowledgeBaseName and DataSourceName as CloudFormation parameters with defaults, and expose them as stack Outputs. This allows 1.Setup.ipynb to retrieve the names via DescribeStacks and look up the manually-created Knowledge Base and Data Source by name, rather than relying on a hardcoded substring match against the stack name.

### 5. Add required IAM permissions for MSK topic management and DescribeStacks

- SageMakerMSKAccessPolicy: Add kafka:CreateTopic, kafka:ListTopics, kafka:DescribeTopic, and kafka-cluster:CreateTopic
- SageMakerCloudFormationAccessPolicy: Add cloudformation:DescribeStacks

## Motivation

- **Simpler setup**: Eliminates the most complex and error-prone step (terminal-based Kafka client installation)
- **Maintained dependencies**: Replaces an archived library with an actively maintained one
- **AWS recommended version**: Uses the AWS recommended Kafka version
- **Better configuration**: Parameterized KB/DS names instead of hardcoded values

## Testing

All changes have been verified by deploying the CloudFormation stack, running all three notebooks sequentially (1.Setup, 2.StreamIngest, 3.Cleanup), and confirming successful stream ingestion into Bedrock Knowledge Base.

I have a pull request ready and will submit it referencing this issue.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.