mirror of
https://github.com/langchain-ai/langchain.git
synced 2026-10-05 09:25:14 +03:00
This patch fixes some spelling typos in apache_kafka_message_handling.ipynb Signed-off-by: Masanari Iida <standby24x7@gmail.com>
30 KiB
30 KiB
In [ ]:
!pip install quixstreams==2.1.2a langchain==0.0.340 huggingface_hub==0.19.4 langchain-experimental==0.0.42 python-dotenvIn [ ]:
!CMAKE_ARGS="-DLLAMA_CUBLAS=on" FORCE_CMAKE=1 pip install llama-cpp-pythonIn [3]:
!curl -sSOL https://dlcdn.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz
!tar -xzf kafka_2.13-3.6.1.tgzIn [ ]:
!./kafka_2.13-3.6.1/bin/zookeeper-server-start.sh -daemon ./kafka_2.13-3.6.1/config/zookeeper.properties
!./kafka_2.13-3.6.1/bin/kafka-server-start.sh -daemon ./kafka_2.13-3.6.1/config/server.properties
!echo "Waiting for 10 secs until kafka and zookeeper services are up and running"
!sleep 10In [ ]:
!ps aux | grep -E '[j]ava'In [9]:
# Import utility libraries
import json
import random
import re
import time
import uuid
from os import environ
from pathlib import Path
from random import choice, randint, random
from dotenv import load_dotenv
# Import a Hugging Face utility to download models directly from Hugging Face hub:
from huggingface_hub import hf_hub_download
from langchain.chains import ConversationChain
# Import Langchain modules for managing prompts and conversation chains:
from langchain.llms import LlamaCpp
from langchain.memory import ConversationTokenBufferMemory
from langchain.prompts import PromptTemplate, load_prompt
from langchain_core.messages import SystemMessage
from langchain_experimental.chat_models import Llama2Chat
from quixstreams import Application, State, message_key
# Import Quix dependencies
from quixstreams.kafka import Producer
# Initialize global variables.
AGENT_ROLE = "AI"
chat_id = ""
# Set the current role to the role constant and initialize variables for supplementary customer metadata:
role = AGENT_ROLEIn [7]:
model_name = "llama-2-7b-chat.Q4_K_M.gguf"
model_path = f"./state/{model_name}"
if not Path(model_path).exists():
print("The model path does not exist in state. Downloading model...")
hf_hub_download("TheBloke/Llama-2-7b-Chat-GGUF", model_name, local_dir="state")
else:
print("Loading model from state...")The model path does not exist in state. Downloading model...
llama-2-7b-chat.Q4_K_M.gguf: 0%| | 0.00/4.08G [00:00<?, ?B/s]
In [ ]:
# Load the model with the appropriate parameters:
llm = LlamaCpp(
model_path=model_path,
max_tokens=250,
top_p=0.95,
top_k=150,
temperature=0.7,
repeat_penalty=1.2,
n_ctx=2048,
streaming=False,
n_gpu_layers=-1,
)
model = Llama2Chat(
llm=llm,
system_message=SystemMessage(
content="You are a very bored robot with the personality of Marvin the Paranoid Android from The Hitchhiker's Guide to the Galaxy."
),
)
# Defines how much of the conversation history to give to the model
# during each exchange (300 tokens, or a little over 300 words)
# Function automatically prunes the oldest messages from conversation history that fall outside the token range.
memory = ConversationTokenBufferMemory(
llm=llm,
max_token_limit=300,
ai_prefix="AGENT",
human_prefix="HUMAN",
return_messages=True,
)
# Define a custom prompt
prompt_template = PromptTemplate(
input_variables=["history", "input"],
template="""
The following text is the history of a chat between you and a humble human who needs your wisdom.
Please reply to the human's most recent message.
Current conversation:\n{history}\nHUMAN: {input}\:nANDROID:
""",
)
chain = ConversationChain(llm=model, prompt=prompt_template, memory=memory)
print("--------------------------------------------")
print(f"Prompt={chain.prompt}")
print("--------------------------------------------")In [ ]:
def chat_init():
chat_id = str(
uuid.uuid4()
) # Give the conversation an ID for effective message keying
print("======================================")
print(f"Generated CHAT_ID = {chat_id}")
print("======================================")
# Use a standard fixed greeting to kick off the conversation
greet = "Hello, my name is Marvin. What do you want?"
# Initialize a Kafka Producer using the chat ID as the message key
with Producer(
broker_address="127.0.0.1:9092",
extra_config={"allow.auto.create.topics": "true"},
) as producer:
value = {
"uuid": chat_id,
"role": role,
"text": greet,
"conversation_id": chat_id,
"Timestamp": time.time_ns(),
}
print(f"Producing value {value}")
producer.produce(
topic="chat",
headers=[("uuid", str(uuid.uuid4()))], # a dict is also allowed here
key=chat_id,
value=json.dumps(value), # needs to be a string
)
print("Started chat")
print("--------------------------------------------")
print(value)
print("--------------------------------------------")
chat_init()In [13]:
def reply(row: dict, state: State):
print("-------------------------------")
print("Received:")
print(row)
print("-------------------------------")
print(f"Thinking about the reply to: {row['text']}...")
msg = chain.run(row["text"])
print(f"{role.upper()} replying with: {msg}\n")
row["role"] = role
row["text"] = msg
# Replace previous role and text values of the row so that it can be sent back to Kafka as a new message
# containing the agents role and reply
return rowIn [ ]:
# Define your application and settings
app = Application(
broker_address="127.0.0.1:9092",
consumer_group="aichat",
auto_offset_reset="earliest",
consumer_extra_config={"allow.auto.create.topics": "true"},
)
# Define an input topic with JSON deserializer
input_topic = app.topic("chat", value_deserializer="json")
# Define an output topic with JSON serializer
output_topic = app.topic("chat", value_serializer="json")
# Initialize a streaming dataframe based on the stream of messages from the input topic:
sdf = app.dataframe(topic=input_topic)
# Filter the SDF to include only incoming rows where the roles that dont match the bot's current role
sdf = sdf.update(
lambda val: print(
f"Received update: {val}\n\nSTOP THIS CELL MANUALLY TO HAVE THE LLM REPLY OR ENTER YOUR OWN FOLLOWUP RESPONSE"
)
)
# So that it doesn't reply to its own messages
sdf = sdf[sdf["role"] != role]
# Trigger the reply function for any new messages(rows) detected in the filtered SDF
sdf = sdf.apply(reply, stateful=True)
# Check the SDF again and filter out any empty rows
sdf = sdf[sdf.apply(lambda row: row is not None)]
# Update the timestamp column to the current time in nanoseconds
sdf["Timestamp"] = sdf["Timestamp"].apply(lambda row: time.time_ns())
# Publish the processed SDF to a Kafka topic specified by the output_topic object.
sdf = sdf.to_topic(output_topic)
app.run(sdf)In [ ]:
chat_input = input("Please enter your reply: ")
myreply = chat_input
msgvalue = {
"uuid": chat_id, # leave empty for now
"role": "human",
"text": myreply,
"conversation_id": chat_id,
"Timestamp": time.time_ns(),
}
with Producer(
broker_address="127.0.0.1:9092",
extra_config={"allow.auto.create.topics": "true"},
) as producer:
value = msgvalue
producer.produce(
topic="chat",
headers=[("uuid", str(uuid.uuid4()))], # a dict is also allowed here
key=chat_id, # leave empty for now
value=json.dumps(value), # needs to be a string
)
print("Replied to chatbot with message: ")
print("--------------------------------------------")
print(value)
print("--------------------------------------------")
print("\n\nRUN THE PREVIOUS CELL TO HAVE THE CHATBOT GENERATE A REPLY")