Upstash Kafka・Redis・Next.jsでリアルタイムチャットアプリを構築する完全ガイド
プロジェクト概要
この記事では、ユーザーがメッセージクライアントやチャットルームを作成でき、過去のメッセージ履歴にもアクセスできるメッセージングアプリケーションを構築します。
本プロジェクトは2つのページで構成されています。1ページ目はクライアント登録専用のページで、一意の名前を持つ複数のクライアントを作成できます。

クライアントのユーザー名をクリックすると、そのユーザーに関連付けられたチャットルーム画面へ移動します。

チャットアプリの仕組み
チャットアプリケーションのロジックは以下の通りです。
ユーザーはインデックスページで一意のユーザー名を持つ複数のクライアントを作成できます。ユーザー名をクリックすると、固有のパスを持つ個別のクライアントが新しいタブで開きます。
各クライアントはWebSocket接続を介してメッセージサーバーに接続されます。クライアント上で新しいメッセージが作成されると、そのクライアントに関連付けられたメッセージサーバーへ送信されます。
メッセージサーバーはメッセージトラフィックを処理します。クライアントがWebSocket接続経由でメッセージを送信すると、サーバーはそのメッセージをKafkaブローカーへ振り分けます。各メッセージサーバーはNode.jsスレッドを実行して受信メッセージを処理し、メッセージがコンシュームされると、既存のWebSocket接続を通じてクライアントへ配信されます。クライアント側での受信メッセージの処理には、react-use-websocketライブラリを使用します。
また、アプリケーションはUpstash Redisを使用してメッセージ履歴を保存します。Kafkaにメッセージがプロデュースされると同時に、Redisデータベースにも永続化されます。新しいクライアントを作成すると、過去のメッセージがUpstash Redisから取得され、チャット画面に表示されます。
以下がアプリケーション全体の構成図です。
注: 本実装ではデモ用に単一のメッセージサーバーを作成していますが、実際にはメッセージ負荷に応じてサーバー数を増やすことが可能です。

デモ
アプリのデモはこちらから確認できます。現在のバージョンはFly.ioにデプロイされています。
はじめに
チャットアプリケーションを構築する手順は以下の通りです。
- Upstash Redisデータベースの作成
- Upstash Kafkaクラスターの作成
- Next.jsアプリケーション(フロントエンド)の作成
- WebSocketメッセージサーバーの作成
- Fly.ioへのアプリケーションのデプロイ
Upstash Redisデータベースの作成
Upstash Consoleにアクセスしてログインし、RedisタブでCreate Databaseボタンをクリックします。

これだけでRedisが使用可能になります!認証情報は後ほどRedisコンソールから取得します。
Upstash Kafkaクラスターの作成
次にKafkaタブに切り替え、Create Clusterボタンをクリックします。クラスタ名を入力して進み、続いてKafkaトピックを作成して確定します。

Next.jsアプリの作成
まず、ターミナルからアプリケーションのルートフォルダを作成して移動します。Next.jsアプリとサーバーの両方をこのフォルダに配置します。
mkdir chat-app
cd chat-app
次に、Next.jsアプリを作成します。
$ npx create-next-app@latest
✔ What is your project named? … next-chat-app
✔ Would you like to use TypeScript? … Yes
✔ Would you like to use ESLint? … Yes
✔ Would you like to use Tailwind CSS? … No
✔ Would you like to use `src/` directory? … No
✔ Would you like to use App Router? (recommended) … No
✔ Would you like to customize the default import alias? … No
認証情報の管理
認証情報を保存するため、.envという名前のファイルを作成します。認証情報を何度もコピー&ペーストする必要はなく、このファイルからインポートするだけです。
まず、.envファイルを作成します。
次に、Redisコンソールに移動し、UPSTASH_REDIS_REST_URLとUPSTASH_REDIS_REST_TOKENの認証情報を.envファイルに貼り付けます。

最後に、Kafkaコンソールに切り替えて、UPSTASH_KAFKA_REST_URL、UPSTASH_KAFKA_REST_USERNAME、UPSTASH_KAFKA_REST_PASSWORDを転記します。

これで、.envファイルは以下のようになります。
.env
UPSTASH_REDIS_REST_URL=...
UPSTASH_REDIS_REST_TOKEN=...
UPSTASH_KAFKA_REST_URL=...
UPSTASH_KAFKA_REST_USERNAME=...
UPSTASH_KAFKA_REST_PASSWORD=...
認証情報の設定が完了したので、アプリケーションの実装に進みましょう。
クライアント登録ページ
インデックスページでは、クライアントの登録・作成操作を行います。ユーザー名が送信されると、新しいクライアントが作成され、「Current Clients」テーブルに追加されます。
pages/index.tsx
import { useState } from "react";
import Link from "next/link";
import { Redis } from "@upstash/redis";
import styles from "@/styles/Home.module.css";
export default function Home() {
const [usernameInput, setUsernameInput] = useState<string>("");
const [usernameList, setUsernameList] = useState<string[]>(Array<string>);
const handleInputChange = (e: React.ChangeEvent<HTMLInputElement>): void => {
const inputValue: string = e.target.value;
setUsernameInput(inputValue);
};
const addUsernameClient = (e: React.FormEvent<HTMLFormElement>): void => {
e.preventDefault();
setUsernameList([...usernameList, usernameInput]);
setUsernameInput("");
};
return (
<div className={styles.container}>
<div className={styles.welcomeSection}>
<h1>Welcome to the demo message app!</h1>
<p>
This application uses Upstash Kafka for message passing, and Upstash
Redis for state management.
<br />
<br />
To get started, create several clients by typing in unique usernames to
the input section below and submitting.
<br />
<br />
The usernames will be added to the list of current clients. Click on a
username to open a new tab with that client's message display.
<br />
<br />
You can have multiple sessions open at once.
</p>
</div>
<form className={styles.formSection} onSubmit={addUsernameClient}>
<input
type="text"
className={styles.formInput}
value={usernameInput}
onChange={handleInputChange}
></input>
<button className={styles.formSubmit} type="submit">
Create the client!
</button>
</form>
<div className={styles.clientListSection}>
<p className={styles.clientListHeader}>Current Clients</p>
<div className={styles.clientList}>
{usernameList.map((username, i) => {
return (
<Link
href={`/user/${username}`}
key={`${i}`}
className={styles.userClient}
target={"_blank"}
>
<p>{username}</p>
</Link>
);
})}
</div>
</div>
</div>
);
}
アプリを再読み込みするたびにチャット履歴をリセットしたい場合は、以下の関数を使用できます。
pages/index.tsx
export async function getServerSideProps() {
const redis = new Redis({
url: process.env.UPSTASH_REDIS_REST_URL,
token: process.env.UPSTASH_REDIS_REST_TOKEN,
});
await redis.del("messagesList");
return {
props: {},
};
}
これでインデックスページの準備が整いました。next-chat-appフォルダ内でnpm run devコマンドを実行すると、インデックスページが起動していることを確認できます。
メッセージクライアントページ
クライアントに対して動的ルーティングを実装するため、/pages/user/[username].tsxというフォルダ構造を作成します。これにより、ユーザー名ごとに個別の動的ルートを生成できます。
以下がクライアントのメインコンポーネントです。このコンポーネントは、メッセージリストやユーザー名などの状態を保持します。useWebSocketフックを使用して、WebSocketからのメッセージイベント、接続イベント、切断イベントを処理します。メッセージイベントが発行されると、メッセージがリストに追加され、MessageDisplayコンポーネントが再レンダリングされます。
/pages/user/[username].tsx
import { useState } from "react";
import { useRouter } from "next/router";
import { Redis } from "@upstash/redis";
import useWebSocket from "react-use-websocket";
import styles from "@/styles/Home.module.css";
type Message = {
id: number;
sender: string;
text: string;
};
export default function MessageApp(props: { messagesData: Message[] }) {
const { messagesData } = props;
const { username } = useRouter().query;
const [inputText, setInputText] = useState<string>("");
const [messageList, setMessageList] = useState<Message[]>(messagesData);
const [messageCounter, setMessageCounter] = useState<number>(0);
const handleMessage = function (message: Message) {
const nextMessages = [...messageList, message];
setMessageList(nextMessages);
};
// WebSocketイベントの処理
const { sendMessage } = useWebSocket("ws://localhost:8080", {
share: true,
filter: () => false,
onOpen: () => {
console.log("WebSocket connection!");
return "connection";
},
onMessage: (message) => {
const data = JSON.parse(message.data);
const { sender, text }: { sender: string; text: string } = data;
const messageData: Message = {
id: messageCounter,
sender: sender,
text: text,
};
setMessageCounter(messageCounter + 1);
handleMessage(messageData);
return message;
},
onClose: () => {
console.log("WebSocket disconnected!");
return "disconnected";
},
});
function handleSendMessage(messageText: string) {
const messageData = {
sender: username,
text: messageText,
};
sendMessage(JSON.stringify(messageData));
}
return (
<div className={styles.Container}>
<MessageDisplay messages={messageList} />
<MessageInput
inputText={inputText}
setInputText={setInputText}
handleSendMessage={handleSendMessage}
/>
</div>
);
}
次に、MessageDisplayコンポーネントとMessageInputコンポーネントです。
/pages/user/[username].tsx
const MessageDisplay = function (props: { messages: Message[] }) {
const { messages } = props;
return (
<div className={styles.messageContainer}>
{messages.map((message) => (
<MessageBubble
key={message.id}
sender={message.sender}
text={message.text}
/>
))}
</div>
);
};
const MessageInput = (props: {
inputText: string;
setInputText: (msg: string) => void;
handleSendMessage: (msg: string) => void;
}) => {
const { inputText, setInputText, handleSendMessage } = props;
const handleInputChange = (
e: React.ChangeEvent<HTMLInputElement>
): void => {
const inputValue: string = e.target.value;
setInputText(inputValue);
};
const handleSubmit = (e: React.FormEvent<HTMLFormElement>): void => {
e.preventDefault();
handleSendMessage(inputText);
if (inputText.trim() !== "") {
setInputText(" ");
}
};
return (
<form className={styles.inputSection} onSubmit={handleSubmit}>
<input
className={styles.inputText}
type="text"
value={inputText}
onChange={handleInputChange}
></input>
<button className={styles.inputSendButton} type="submit">
Send
</button>
</form>
);
};
const MessageBubble = (props: {
sender: string;
text: string;
key: number;
}) => {
const { sender, text } = props;
const { username } = useRouter().query;
const isSender = sender === username;
const senderClass = isSender ? "sender" : "receiver";
return (
<div className={`${styles["messageBubble"]} ${styles[senderClass]}`}>
<div className={styles.messageSender}>
{isSender ? "You" : sender}
</div>
<div className={styles.messageText}>{text}</div>
</div>
);
};
クライアントにチャット履歴を提供するために、getServerSideProps()関数を使用します。
/pages/user/[username].tsx
export async function getServerSideProps() {
const redis = new Redis({
url: process.env.UPSTASH_REDIS_REST_URL,
token: process.env.UPSTASH_REDIS_REST_TOKEN,
});
const messagesData = (await redis.lrange("messagesList", 0, -1)).reverse();
return {
props: {
messagesData,
},
};
}
これでNext.jsアプリが動作するようになりました。ページを更新してクライアントを作成し、いずれかのクライアントに移動してみてください。クライアントページが表示されるはずです。ただし、メッセージフローを処理するためのメッセージサーバーがまだ必要です。
メッセージサーバーの作成
サーバーの構造は非常にシンプルです。Node.js、wsライブラリ、Upstash Kafkaを使用して実装します。まず、chat-appフォルダ内にserverフォルダを作成します。
mkdir server
cd server
serverフォルダ内で、必要なパッケージをインストールし、設定ファイルを生成します。
npm install typescript ws tsc @upstash/kafka @types/ws
tsc --init
次に、/server/message_server.tsファイル内にWebSocket、Kafkaプロデューサー、Kafkaコンシューマーのクライアントを作成します。
/server/message_server.ts
import * as http from "http";
import { Kafka } from "@upstash/kafka";
import { Redis } from "@upstash/redis";
import { WebSocket } from "ws";
const server = http.createServer();
const wss = new WebSocket.Server({ server });
server.listen(8080, () => {
console.log("Server is running on port 8080");
});
const kafka = new Kafka({
url: process.env.UPSTASH_KAFKA_REST_URL,
username: process.env.UPSTASH_KAFKA_REST_USERNAME,
password: process.env.UPSTASH_KAFKA_REST_PASSWORD,
});
const redis = new Redis({
url: process.env.UPSTASH_REDIS_REST_URL,
token: process.env.UPSTASH_REDIS_REST_TOKEN,
});
const consumer = kafka.consumer();
const producer = kafka.producer();
const clients = new Set<WebSocket>();
WebSocketとやり取りするために、connectionイベントとmessageイベントを作成します。
/server/message_server.ts
wss.on("connection", async (connection, req) => {
clients.add(connection);
console.log(`New client connected!`);
connection.on("message", async (message) => {
const jsonMessage = message.toString();
console.log("Received message:", JSON.parse(jsonMessage));
producer.produce("chat", jsonMessage);
});
connection.on("close", () => {
console.log(`Client disconnected:`);
clients.delete(connection);
});
});
最後に、定義済みの間隔でメッセージをコンシュームするスレッドを作成して実行します。
/server/message_server.ts
async function run() {
while (true) {
const messages = await consumer.consume({
consumerGroupId: "group_1",
instanceId: "instance_1",
topics: ["chat"],
autoOffsetReset: "earliest",
});
if (messages.length != 0) {
for (let i = 0; i < messages.length; i++) {
await redis.lpush("messagesList", messages[i].value);
console.log(`Message sending: ${messages[i].value}`);
clients.forEach((connection: WebSocket) => {
connection.send(messages[i].value);
});
}
}
console.log("Run!");
await new Promise((resolve) => setTimeout(resolve, 1000));
}
}
すべての準備が整いました。この時点でアプリは問題なく動作するはずです。ローカルでメッセージサーバーを起動し、クライアントページを更新すると、クライアント間でメッセージがやり取りされている様子を確認できます。以下のコマンドでTSファイルをトランスパイルし、localhost:8000でサーバーを起動できます。
tsc message_server.ts
node message_server.js
デプロイ
デプロイにはFly.ioを使用します。まだアカウントをお持ちでない場合は、作業を始める前にアカウントを作成してください。
メッセージサーバーのデプロイ
serverフォルダに移動し、flyctl CLIツールをインストールして、シェル経由で認証を行います。
npm install flyctl
flyctl auth login
設定ファイルを作成するためにflyctl initを実行します。これによりfly.tomlが生成されます。fly.tomlを開き、WebSocket接続の設定として以下の行を挿入してください。
fly.toml
[[services]]
internal_port = 8080
protocol = "tcp"
[services.concurrency]
hard_limit = 25
soft_limit = 20
[[services.ports]]
handlers = ["http"]
port = "80"
[[services.ports]]
handlers = ["tls", "http"]
port = "443"
[[services.tcp_checks]]
interval = 10000
timeout = 2000
サーバー側の最後のステップです。flyctl deployを実行すれば準備完了です!デプロイプロセスが完了すると、flyctlがサーバーのエンドポイントを表示します。そのエンドポイントをコピーしておきましょう。今回の例では、エンドポイントはmessage-server.fly.devです。
Next.jsアプリケーションのデプロイ
Next.jsアプリケーションをデプロイする前に、メッセージサーバーのデプロイ済みエンドポイントを組み込む必要があります。pages/user/[username].tsxファイル内のWebSocket URLを、ws://localhost:8080からflyctlで取得したエンドポイントに置き換え、先頭にwss://プレフィックスを付けてください。今回の例ではwss://message-server.fly.devとなります。
次に、next-chat-appフォルダ内で、serverフォルダと同じコマンドを実行します。今回はfly.tomlファイルの編集は不要なので、そのステップは省略できます。
flyctl init
flyctl deploy
これで完成です!flyctl openコマンドを実行すると、デプロイされたプロジェクトにアクセスできます。
まとめと今後の改善案
最後までお読みいただきありがとうございました!
プロジェクトのGitHubリポジトリはこちらからご覧いただけます。
さらに開発を進めたい方向けに、いくつか改善案をご紹介します。
-
現状では、ページを再読み込みするたびに、Upstash Redisに保存されていたすべてのメッセージが削除されます。この挙動は
pages/index.tsxファイル内のgetServerSideProps関数によって制御されています。しかし、ユーザーがページを再読み込みすると、チャットルームに参加している全員のチャット履歴が削除されるという重大な問題が発生します。
この問題を解決するには、メッセージが送信されるたびにチャット履歴のTTL(有効期限)を延長する仕組みを実装することが推奨されます。この改善により、ページを再読み込みしてもチャット履歴が保持され、引き続き参照できるようになります。 -
複数チャットルーム機能の実装も可能です。これを実現するには、チャットルームごとに一意の名前を持つ複数のKafkaトピックを作成する方法があります。あるいは、適切なデータ構造を使ってメッセージサーバー側で処理することも可能です。
-
複数のメッセージサーバーとロードバランサーを実装すれば、より優れたシステム設計のベストプラクティスを適用できます。
ご質問がある場合は、fahreddin@upstash.com までお気軽にお問い合わせください。
-
加速するデジタルトランスフォーメーション──あなたのデータレイヤーは準備できていますか?
この10年間、デジタルトランスフォーメーション(DX)については数え切れないほど語られてきました。Redisでは、変革を可能にする一連の技術レイヤー(クラウド、マイクロサービス、コンテナ、NoSQLデータベース)と、企業が変革の旅をどれだけ速く、成功裏に推進できるかとの間に、直接的な関係があると捉えています。 しかし、この関係をどうやって実行可能なインサイトへと定量化すればよいのでしょうか? そこで登場するのが「Digital Transformation Index(DTI:デジタルトランスフォーメーション指数)」です。まず、上述の4つの技術の採用状況を組み合わせて、企業が変革の道程のどこ
-
Redis PSUBSCRIBEコマンド徹底解説 – Pub/Subで複数パターンを購読する方法
このチュートリアルでは、Redisのメッセージブローカー(Pub/Sub)システムにおいて、redis-cliを使って複数のパターンを同時にサブスクライブ(購読)する方法を解説します。 PSUBSCRIBEコマンドとは PSUBSCRIBEコマンドは、クライアントを1つ以上のパターンに登録し、指定されたパターンに名前が一致するチャンネルへパブリッシュされたすべてのメッセージを受信するためのコマンドです。パターンはglobスタイルで指定します。 SUBSCRIBEコマンドと同様に、クライアントがPSUBSCRIBEを実行すると「サブスクライブ状態」に入り、購読中のパターンを待ち受けます。他のクラ