You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Memgraph连接Kafka Stream Consumer时Broker传输失败问题求助

问题描述

我通过官方docker-compose.yml启动了Memgraph与Confluent Kafka平台,创建了myTopic作为数据传输的Kafka Topic。按照文档实现了名为my_stream的消费者并绑定自定义转换模块,流配置如图所示。添加转换模块并尝试连接流时,抛出错误:

Transformation Failed
Failed to initialize Kafka consumer my_stream : Local: Broker transport failure

运行docker-compose.yml的终端中有更详细的错误信息(如图)。Confluent集群概览显示Broker运行正常,用文档中的movielens流测试同一转换模块时,流能正常连接(仅转换阶段报错),因此暂排除转换模块问题。自定义转换模块代码如下:

import mgp 
import json

@mgp.transformation
def transfer(messages: mgp.Messages) -> mgp.Record(query=str, parameters=mgp.Nullable[mgp.Map]):
  result_queries = []

  for i in range(messages.total_messages()):
    message = messages.message_at(i)
    transfer_dict = json.loads(message.payload().decode('utf8'))
    result_queries.append(
      mgp.Record(
        query=(
          """
        MERGE (fromAccount:Account {accountId: $fromId, createTime: $fCreateTime, isBlocked: $fIsBlocked, accountType: $fAccountType, nickname: $fNickname, phonenum: $fPhonenum, 
        email: $fEmail, freqLoginType: $fFreqLoginType, lastLoginTime: $fLastLoginTime, accountLevel: $fAccountLevel})
        MERGE (toAccount:Account {accountId: $toId, createTime: $tCreateTime, isBlocked: $tIsBlocked, accountType: $tAccountType, nickname: $tNickname, phonenum: $tPhonenum, 
        email: $tEmail, freqLoginType: $tFreqLoginType, lastLoginTime: $tLastLoginTime, accountLevel: $tAccountLevel})
        MERGE (fromAccount)-[:TRANSFER {amount: $amount, createTime: $eCreateTime, orderNum: $orderNum, comment: $comment, payType: $payType, goodsType: $goodsType, 
        insertTime: $insertTime, exportTime: timestamp()}]->(toAccount)
        """),
        parameters={
          "fromId": transfer_dict["fromAccount"]["fromId"], "fCreateTime": transfer_dict["fromAccount"]["createTime"], "fIsBlocked": transfer_dict["fromAccount"]["isBlocked"],
          "fAccountType": transfer_dict["fromAccount"]["accountType"], "fNickname": transfer_dict["fromAccount"]["nickname"], "fPhonenum": transfer_dict["fromAccount"]["phonenum"], 
          "fEmail": transfer_dict["fromAccount"]["email"], "fFreqLoginType": transfer_dict["fromAccount"]["freqLoginType"], "fLastLoginTime": transfer_dict["fromAccount"]["lastLoginTime"], 
          "fAccountLevel": transfer_dict["fromAccount"]["accountLevel"],
          "toId": transfer_dict["toAccount"]["toId"], "tCreateTime": transfer_dict["toAccount"]["createTime"], "tIsBlocked": transfer_dict["toAccount"]["isBlocked"],
          "tAccountType": transfer_dict["toAccount"]["accountType"], "tNickname": transfer_dict["toAccount"]["nickname"], "tPhonenum": transfer_dict["toAccount"]["phonenum"],
          "tEmail": transfer_dict["toAccount"]["email"], "tFreqLoginType": transfer_dict["toAccount"]["freqLoginType"], "tLastLoginTime": transfer_dict["toAccount"]["lastLoginTime"],
          "tAccountLevel": transfer_dict["toAccount"]["accountLevel"],
          "amount": transfer_dict["amount"], "eCreateTime": transfer_dict["createTime"], "orderNum": transfer_dict["orderNum"], "comment": transfer_dict["comment"], 
          "payType": transfer_dict["payType"], "goodsType": transfer_dict["goodsType"], "insertTime": transfer_dict["insertTime"]
        }
      )
    )
    return result_queries
原因分析
  1. Docker网络连通性问题:Memgraph容器与Kafka容器可能不在同一网络,或者流配置中指定的Kafka Broker地址无法被Memgraph容器解析。movielens流能正常连接,说明转换模块无问题,故障点集中在my_stream的Kafka连接配置上。
  2. Kafka Topic状态异常:myTopic可能存在创建未生效、分区/副本配置不符合要求的情况,导致Memgraph消费者无法正常连接。
  3. 安全配置不匹配:若Confluent Kafka启用了SASL/SSL认证,而my_stream未配置对应的认证参数,会引发传输层失败。
解决方法

1. 修复Docker网络配置

  • 确认Memgraph与Confluent Kafka容器处于同一自定义网络:检查docker-compose.yml,若未配置共享网络,添加如下配置:
    networks:
      kafka-network:
        driver: bridge
    
    并在两个服务的配置中添加networks: [kafka-network]字段。
  • 流配置中使用Kafka容器的服务名作为Broker地址(如kafka:9092),而非宿主机IP或localhost——容器内部无法直接解析宿主机的localhost。

2. 验证Kafka Topic状态

  • 进入Confluent Kafka容器,执行命令检查myTopic的存在与状态:
    kafka-topics --list --bootstrap-server kafka:9092
    kafka-topics --describe --topic myTopic --bootstrap-server kafka:9092
    
    若Topic不存在,重新创建:
    kafka-topics --create --topic myTopic --bootstrap-server kafka:9092 --partitions 1 --replication-factor 1
    

3. 匹配安全认证配置

若Confluent Kafka启用了认证,在Memgraph流配置中添加对应参数,示例如下:

CREATE STREAM my_stream
TOPICS myTopic
TRANSFORM transfer
BOOTSTRAP_SERVERS 'kafka:9092'
SECURITY_PROTOCOL 'SASL_PLAINTEXT'
SASL_MECHANISM 'PLAIN'
SASL_USERNAME 'your-username'
SASL_PASSWORD 'your-password';

4. 修复转换模块的逻辑错误

你的转换模块中return result_queries缩进错误,位于for循环内部,导致仅处理第一条消息就返回。修正后代码如下:

import mgp 
import json

@mgp.transformation
def transfer(messages: mgp.Messages) -> mgp.Record(query=str, parameters=mgp.Nullable[mgp.Map]):
  result_queries = []

  for i in range(messages.total_messages()):
    message = messages.message_at(i)
    transfer_dict = json.loads(message.payload().decode('utf8'))
    result_queries.append(
      mgp.Record(
        query=(
          """
        MERGE (fromAccount:Account {accountId: $fromId, createTime: $fCreateTime, isBlocked: $fIsBlocked, accountType: $fAccountType, nickname: $fNickname, phonenum: $fPhonenum, 
        email: $fEmail, freqLoginType: $fFreqLoginType, lastLoginTime: $fLastLoginTime, accountLevel: $fAccountLevel})
        MERGE (toAccount:Account {accountId: $toId, createTime: $tCreateTime, isBlocked: $tIsBlocked, accountType: $tAccountType, nickname: $tNickname, phonenum: $tPhonenum, 
        email: $tEmail, freqLoginType: $tFreqLoginType, lastLoginTime: $tLastLoginTime, accountLevel: $tAccountLevel})
        MERGE (fromAccount)-[:TRANSFER {amount: $amount, createTime: $eCreateTime, orderNum: $orderNum, comment: $comment, payType: $payType, goodsType: $goodsType, 
        insertTime: $insertTime, exportTime: timestamp()}]->(toAccount)
        """),
        parameters={
          "fromId": transfer_dict["fromAccount"]["fromId"], "fCreateTime": transfer_dict["fromAccount"]["createTime"], "fIsBlocked": transfer_dict["fromAccount"]["isBlocked"],
          "fAccountType": transfer_dict["fromAccount"]["accountType"], "fNickname": transfer_dict["fromAccount"]["nickname"], "fPhonenum": transfer_dict["fromAccount"]["phonenum"], 
          "fEmail": transfer_dict["fromAccount"]["email"], "fFreqLoginType": transfer_dict["fromAccount"]["freqLoginType"], "fLastLoginTime": transfer_dict["fromAccount"]["lastLoginTime"], 
          "fAccountLevel": transfer_dict["fromAccount"]["accountLevel"],
          "toId": transfer_dict["toAccount"]["toId"], "tCreateTime": transfer_dict["toAccount"]["createTime"], "tIsBlocked": transfer_dict["toAccount"]["isBlocked"],
          "tAccountType": transfer_dict["toAccount"]["accountType"], "tNickname": transfer_dict["toAccount"]["nickname"], "tPhonenum": transfer_dict["toAccount"]["phonenum"],
          "tEmail": transfer_dict["toAccount"]["email"], "tFreqLoginType": transfer_dict["toAccount"]["freqLoginType"], "tLastLoginTime": transfer_dict["toAccount"]["lastLoginTime"],
          "tAccountLevel": transfer_dict["toAccount"]["accountLevel"],
          "amount": transfer_dict["amount"], "eCreateTime": transfer_dict["createTime"], "orderNum": transfer_dict["orderNum"], "comment": transfer_dict["comment"], 
          "payType": transfer_dict["payType"], "goodsType": transfer_dict["goodsType"], "insertTime": transfer_dict["insertTime"]
        }
      )
    )
  # 将return移至循环外部
  return result_queries

内容的提问来源于stack exchange,提问作者Max

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 09:44:52