Redis Pub/Sub入門:Node.jsとTypeScriptで実装する

スポンサーリンク

Redis Pub/Sub入門:Node.jsとTypeScriptで実装する

Pub/Subとは

Pub/Sub(Publish/Subscribe)はメッセージの送信者(Publisher)と受信者(Subscriber)を直接結びつけないメッセージングパターンです。PublisherはチャンネルにメッセージをPublishし、そのチャンネルをSubscribeしているすべての受信者に配信されます。

Publisher → [channel: notifications] → Subscriber A
                                      → Subscriber B
                                      → Subscriber C

RedisはPub/Subをネイティブにサポートしており、シンプルなAPIで実装できます。


準備

Redisを起動する(Docker)

docker run -d --name redis -p 6379:6379 redis:latest

Node.jsプロジェクトのセットアップ

mkdir redis-pubsub && cd redis-pubsub
npm init -y
npm install ioredis
npm install -D typescript @types/node ts-node
npx tsc --init

基本実装

Subscriber(受信側)

// subscriber.ts
import Redis from 'ioredis'

const subscriber = new Redis({ host: 'localhost', port: 6379 })

// チャンネルをSubscribe
await subscriber.subscribe('notifications')

subscriber.on('message', (channel: string, message: string) => {
  console.log(`[${channel}] ${message}`)
})

console.log('Subscribing to notifications...')

Publisher(送信側)

// publisher.ts
import Redis from 'ioredis'

const publisher = new Redis({ host: 'localhost', port: 6379 })

// チャンネルにメッセージをPublish
await publisher.publish('notifications', 'Hello, Subscribers!')
await publisher.publish('notifications', JSON.stringify({ type: 'alert', body: '新しい通知があります' }))

await publisher.quit()

重要: subscribe を呼んだ接続はSubscribeモードに入り、PUBLISH などの通常コマンドを実行できません。PublisherとSubscriberで接続を分ける必要があります。


TypeScriptで型安全にする

メッセージの型をジェネリクスで管理します。

// types.ts
export type NotificationMessage = {
  type: 'info' | 'warning' | 'error'
  body: string
  userId?: number
}

export type OrderMessage = {
  orderId: number
  status: 'created' | 'shipped' | 'delivered'
}
// subscriber.ts
import Redis from 'ioredis'
import type { NotificationMessage, OrderMessage } from './types'

const subscriber = new Redis()

await subscriber.subscribe('notifications', 'orders')

subscriber.on('message', (channel: string, message: string) => {
  if (channel === 'notifications') {
    const data = JSON.parse(message) as NotificationMessage
    console.log(`通知 [${data.type}]: ${data.body}`)
  }

  if (channel === 'orders') {
    const data = JSON.parse(message) as OrderMessage
    console.log(`注文 #${data.orderId}: ${data.status}`)
  }
})
// publisher.ts
import Redis from 'ioredis'
import type { NotificationMessage, OrderMessage } from './types'

const publisher = new Redis()

const notify = (msg: NotificationMessage) =>
  publisher.publish('notifications', JSON.stringify(msg))

const updateOrder = (msg: OrderMessage) =>
  publisher.publish('orders', JSON.stringify(msg))

await notify({ type: 'info', body: 'デプロイが完了しました', userId: 42 })
await updateOrder({ orderId: 1001, status: 'shipped' })

await publisher.quit()

チャンネルパターン(psubscribe)

psubscribe を使うとワイルドカードで複数チャンネルを一括購読できます。

// パターンで複数チャンネルを購読
await subscriber.psubscribe('user:*')

subscriber.on('pmessage', (pattern: string, channel: string, message: string) => {
  // pattern: 'user:*'
  // channel: 'user:42', 'user:100' など
  console.log(`[${channel}] ${message}`)
})
// publisher側でチャンネルを動的に指定
await publisher.publish('user:42', JSON.stringify({ event: 'login' }))
await publisher.publish('user:100', JSON.stringify({ event: 'purchase' }))

ユーザーIDやテナントIDをチャンネル名に埋め込んで、対象を絞った通知を実装できます。


実用例:非同期通知システム

APIサーバーがイベントをPublishし、通知ワーカーがSubscribeして処理する例です。

// api-server.ts(PublisherとしてRedisに接続)
import Redis from 'ioredis'
import type { NotificationMessage } from './types'

const redis = new Redis()

// 注文完了時にイベントを発行
export async function onOrderCompleted(userId: number, orderId: number) {
  const message: NotificationMessage = {
    type: 'info',
    body: `注文 #${orderId} が確定しました`,
    userId,
  }
  await redis.publish('notifications', JSON.stringify(message))
}
// notification-worker.ts(Subscriberとして常駐)
import Redis from 'ioredis'
import type { NotificationMessage } from './types'

const subscriber = new Redis()

await subscriber.subscribe('notifications')

subscriber.on('message', async (_channel: string, message: string) => {
  const data = JSON.parse(message) as NotificationMessage
  if (data.userId) {
    // メール送信・プッシュ通知などの処理
    await sendPushNotification(data.userId, data.body)
  }
})

async function sendPushNotification(userId: number, body: string) {
  console.log(`[Push] userId: ${userId}${body}`)
}

注意点

注意事項 内容
メッセージは永続化されない Subscribeしていないタイミングのメッセージは消える
配信保証なし ネットワーク障害時にメッセージが失われる可能性がある
Publisher/Subscriberで接続を分ける subscribe後の接続ではPUBLISHなど通常コマンドが使えない
再接続時の取りこぼし ioredisのautoResubscribe: true(デフォルト)で再接続後に自動再購読される

メッセージの永続化・順序保証・少なくとも1回の配信が必要な場合はRedis Streamsまたは BullMQ の使用を検討します。


まとめ

操作 コマンド / メソッド
チャンネルを購読 subscriber.subscribe('channel')
パターンで購読 subscriber.psubscribe('channel:*')
メッセージを受け取る subscriber.on('message', callback)
パターンで受け取る subscriber.on('pmessage', callback)
メッセージを送信 publisher.publish('channel', message)
  • PublisherとSubscriberは必ず別の接続を使う
  • メッセージはJSONで渡してTypeScriptの型でキャストする
  • psubscribe でチャンネルパターンを使うとユーザー単位・テナント単位の通知が実装しやすい
  • 永続化・再送が必要な場合はRedis StreamsやBullMQを検討する

ジョブキューとしてRedisを活用したい場合は「BullMQ入門:Node.js(TypeScript)でジョブキューを実装する」も参照してください。