Skip to main content التثبيت باستخدام Codex أو Claude انسخ هذا Prompt والصقه في Codex أو Claude أو مساعد آخر ليراجع صفحة Skill ويثبّتها لك.
يتجاوز الأمر المباشر Prompt المخصّص للمراجعة. افحص المصدر قبل تشغيله.
npx skills add https://github.com/pluginagentmarketplace/custom-plugin-server-side-game-dev --skill message-queuesيبقى الأمر في سطر واحد. مرّر أفقيًا لمراجعته كاملًا قبل النسخ.
تفضّل نسخة محلية؟ نزّل الملفات المتاحة حاليًا لدى SkillsMP.

pluginagentmarketplace
pluginagentmarketplace/custom-plugin-server-side-game-dev
فتح مستودع GitHubالمهن ذات الصلةSOC
استنادا إلى تصنيف SOC المهني
| name | message-queues |
| description | Message queue systems for game servers including Kafka, RabbitMQ, and actor models |
| sasmp_version | 1.3.0 |
| version | 2.0.0 |
| bonded_agent | 01-game-server-architect |
| bond_type | SECONDARY_BOND |
| parameters | {"required":["queue_system"],"optional":["batch_size","ack_mode"],"validation":{"queue_system":{"type":"string","enum":["kafka","rabbitmq","redis_pubsub","nats","sqs"]},"batch_size":{"type":"integer","min":1,"max":1000,"default":100},"ack_mode":{"type":"string","enum":["auto","manual","batch"],"default":"manual"}}} |
| retry_config | {"max_attempts":5,"backoff":"exponential","initial_delay_ms":100,"max_delay_ms":30000,"retryable_errors":["CONNECTION_LOST","BROKER_UNAVAILABLE"]} |
| observability | {"logging":{"level":"info","fields":["queue","topic","partition","offset"]},"metrics":[{"name":"messages_published_total","type":"counter"},{"name":"messages_consumed_total","type":"counter"},{"name":"consumer_lag","type":"gauge"},{"name":"processing_duration_ms","type":"histogram"}]} |
Message Queues for Game Servers
Implement asynchronous messaging for scalable game server architecture.
Queue Systems Comparison
| System | Throughput | Latency | Ordering | Use Case |
|---|
| Kafka | Very High | Medium | Partition | Analytics, events |
| RabbitMQ | High | Low | Queue | Game events |
| Redis Pub/Sub | Very High | Very Low | None | Real-time updates |
| NATS | Very High | Ultra Low | Stream | Game state sync |
| SQS | High | Medium | FIFO option | Cloud native |
Apache Kafka for Game Analytics
Properties producerProps = new Properties();
producerProps.put("bootstrap.servers", "kafka:9092");
producerProps.put("key.serializer", StringSerializer.class.getName());
producerProps.put("value.serializer", JsonSerializer.class.getName());
producerProps.put("acks", "all");
producerProps.put("retries", 3);
producerProps.put("linger.ms", 5);
producerProps.put("batch.size", 16384);
KafkaProducer<String, GameEvent> producer = new KafkaProducer<>(producerProps);
public {
ProducerRecord<String, GameEvent> record = <>(
,
event.getPlayerId(),
event
);
producer.send(record, (metadata, exception) -> {
(exception != ) {
log.error(, exception);
}
});
}
();
consumerProps.put(, );
consumerProps.put(, );
consumerProps.put(, );
consumerProps.put(, );
KafkaConsumer<String, GameEvent> consumer = <>(consumerProps);
consumer.subscribe(List.of());
(running) {
ConsumerRecords<String, GameEvent> records = consumer.poll(Duration.ofMillis());
(ConsumerRecord<String, GameEvent> record : records) {
processEvent(record.value());
}
consumer.commitSync();
}
void
publishEvent
(GameEvent event)
new
ProducerRecord
"game-events"
if
null
"Failed to publish event"
Properties
consumerProps
=
new
Properties
"bootstrap.servers"
"kafka:9092"
"group.id"
"analytics-consumer"
"auto.offset.reset"
"earliest"
"enable.auto.commit"
false
new
KafkaConsumer
"game-events"
while
100
for
RabbitMQ for Game Commands
func connectRabbitMQ() (*amqp.Connection, error) {
var conn *amqp.Connection
var err error
for i := 0; i < 5; i++ {
conn, err = amqp.Dial("amqp://guest:guest@localhost:5672/")
if err == nil {
return conn, nil
}
time.Sleep(time.Second * time.Duration(1<<i))
}
return nil, fmt.Errorf("failed to connect after retries: %w", err)
}
func publishMatchEvent(ch *amqp.Channel, event MatchEvent) error {
body, err := json.Marshal(event)
if err != nil {
return err
}
return ch.Publish(
"game-exchange",
"match.created",
false,
false,
amqp.Publishing{
ContentType: "application/json",
Body: body,
DeliveryMode: amqp.Persistent,
MessageId: uuid.New().String(),
Timestamp: time.Now(),
},
)
}
func consumeMatchEvents(ch *amqp.Channel) error {
msgs, err := ch.Consume(
"match-events",
"",
false,
false,
false,
false,
nil,
)
if err != nil {
return err
}
for msg := range msgs {
var event MatchEvent
if err := json.Unmarshal(msg.Body, &event); err != nil {
msg.Nack(false, false)
continue
}
if err := processMatchEvent(event); err != nil {
msg.Nack(false, true)
continue
}
msg.Ack(false)
}
return nil
}
Redis Pub/Sub for Real-Time
import redis
import json
from concurrent.futures import ThreadPoolExecutor
class GameStatePublisher:
def __init__(self):
self.redis = redis.Redis(host='localhost', port=6379)
def broadcast_state(self, game_id: str, state: dict):
channel = f"game:{game_id}"
self.redis.publish(channel, json.dumps(state))
def broadcast_chat(self, game_id: str, message: dict):
channel = f"chat:{game_id}"
self.redis.publish(channel, json.dumps(message))
class GameStateSubscriber:
def __init__(self, game_id: str, callback):
self.redis = redis.Redis(host='localhost', port=6379)
self.pubsub = self.redis.pubsub()
self.callback = callback
self.game_id = game_id
def subscribe(self):
self.pubsub.subscribe(f"game:{self.game_id}")
for message in self.pubsub.listen():
if message['type'] == 'message':
data = json.loads(message['data'])
self.callback(data)
def unsubscribe(self):
self.pubsub.unsubscribe()
self.pubsub.close()
pool = redis.ConnectionPool(host='localhost', port=6379, max_connections=100)
redis_client = redis.Redis(connection_pool=pool)
NATS for Low-Latency Messaging
func setupNATS() (*nats.Conn, nats.JetStreamContext, error) {
nc, err := nats.Connect("nats://localhost:4222",
nats.RetryOnFailedConnect(true),
nats.MaxReconnects(10),
nats.ReconnectWait(time.Second),
)
if err != nil {
return nil, nil, err
}
js, err := nc.JetStream()
if err != nil {
return nil, nil, err
}
_, err = js.AddStream(&nats.StreamConfig{
Name: "GAME_EVENTS",
Subjects: []string{"game.>"},
Retention: nats.LimitsPolicy,
MaxAge: time.Hour * 24,
Storage: nats.FileStorage,
Replicas: 3,
})
return nc, js, err
}
func publishGameEvent(js nats.JetStreamContext, event GameEvent) error {
data, _ := json.Marshal(event)
ack, err := js.Publish(
fmt.Sprintf("game.%s.%s", event.GameID, event.Type),
data,
)
if err != nil {
return err
}
log.Printf("Published: seq=%d", ack.Sequence)
return nil
}
func consumeGameEvents(js nats.JetStreamContext) error {
sub, err := js.Subscribe("game.>",
func(msg *nats.Msg) {
var event GameEvent
json.Unmarshal(msg.Data, &event)
processEvent(event)
msg.Ack()
},
nats.Durable("game-processor"),
nats.ManualAck(),
nats.AckWait(time.Second*30),
)
if err != nil {
return err
}
defer sub.Unsubscribe()
<-make(chan struct{})
return nil
}
Actor Model (Akka/Orleans)
public interface IPlayerGrain : IGrainWithStringKey
{
Task<PlayerState> GetState();
Task<bool> TakeDamage(int amount, string sourceId);
Task<bool> ApplyBuff(Buff buff);
}
public class PlayerGrain : Grain, IPlayerGrain
{
private readonly IPersistentState<PlayerState> _state;
private readonly ILogger<PlayerGrain> _logger;
public PlayerGrain(
[PersistentState("player", "gameStore")] IPersistentState<PlayerState> state,
ILogger<PlayerGrain> logger)
{
_state = state;
_logger = logger;
}
public Task<PlayerState> GetState() => Task.FromResult(_state.State);
public async Task<bool> TakeDamage(int amount, string sourceId)
{
_state.State.Health -= amount;
if (_state.State.Health <= 0)
{
var gameGrain = GrainFactory.GetGrain<IGameGrain>(_state.State.GameId);
await gameGrain.OnPlayerDeath(this.GetPrimaryKeyString(), sourceId);
}
await _state.WriteStateAsync();
return _state.State.Health > 0;
}
}
var host = new HostBuilder()
.UseOrleans(siloBuilder =>
{
siloBuilder
.UseLocalhostClustering()
.AddRedisGrainStorage("gameStore", options =>
{
options.ConnectionString = "localhost:6379";
})
.ConfigureLogging(logging => logging.AddConsole());
})
.Build();
Use Case Mapping
| Use Case | Recommended | Reason |
|---|
| Cross-server chat | RabbitMQ | Reliable delivery |
| Analytics pipeline | Kafka | High throughput, replay |
| Real-time state | Redis Pub/Sub | Ultra-low latency |
| Distributed game state | Orleans/Akka | Location transparency |
| Match results | Kafka | Ordered, durable |
| Notifications | NATS | Simple, fast |
Troubleshooting
Common Failure Modes
| Error | Root Cause | Solution |
|---|
| Consumer lag | Slow processing | Scale consumers |
| Message loss | Auto-ack before process | Manual ack |
| Duplicate processing | At-least-once | Idempotent handlers |
| Broker unavailable | Single point | Cluster mode |
Debug Checklist
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group analytics-consumer
rabbitmqctl list_queues name messages consumers
redis-cli PUBSUB CHANNELS "game:*"
nats stream info GAME_EVENTS
Unit Test Template
func TestMessagePublishing(t *testing.T) {
container := setupRabbitMQContainer(t)
defer container.Terminate(context.Background())
conn, _ := amqp.Dial(container.URI)
ch, _ := conn.Channel()
event := MatchEvent{
MatchID: "match-123",
EventType: "created",
}
err := publishMatchEvent(ch, event)
require.NoError(t, err)
msgs, _ := ch.Consume("match-events", "", true, false, false, false, nil)
select {
case msg := <-msgs:
var received MatchEvent
json.Unmarshal(msg.Body, &received)
assert.Equal(t, event.MatchID, received.MatchID)
case <-time.After(time.Second * 5):
t.Fatal("timeout waiting for message")
}
}
Resources
assets/ - Queue configurations
references/ - Messaging patterns