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

如何将Kafka采集的数据存储到HDFS?及存储时出现域名解析错误的问题求助

如何将Kafka采集的数据存储到HDFS?及存储时出现域名解析错误的问题求助

我正在做一个学校项目,流程是用Kafka采集数据,将数据存储到HDFS中,之后用Spark进行分析。目前Kafka部分已经能够正常采集数据,但在尝试将数据写入HDFS时遇到了域名解析错误,具体报错信息如下:

Erreur lors de l'enregistrement dans HDFS : HTTPConnectionPool(host='hadoop-worker2', port=9864): Max retries exceeded with url: /webhdfs/v1/data/openAQ/openAQ_data.json?op=APPEND&user.name=hdfs&namenoderpcaddress=hadoop-master:9000&user.name=hdfs (Caused by NameResolutionError("<urllib3.connection.HTTPConnection object at 0x00000276542F1700>: Failed to resolve 'hadoop-worker2' ([Errno 11001] getaddrinfo failed)"))

以下是我实现的代码:

from hdfs import InsecureClient


#### Déclaration des variables

KAFKA_BROKER = "localhost:9092"
KAFKA_TOPIC = "BigDataProjet"

# Configuration du client HDFS
HDFS_URL = "http://localhost:9870"  # URL du Namenode
HDFS_PATH = "/data/openAQ"
HDFS_FILE_PATH = f"{HDFS_PATH}/openAQ_data.json"

client = InsecureClient(HDFS_URL, user="hdfs")


# Initialisation du consumer kafka

consumer = KafkaConsumer(
    KAFKA_TOPIC,
    bootstrap_servers=KAFKA_BROKER,
    auto_offset_reset="earliest",
    enable_auto_commit=True
)


def sauvegarder_hdfs(data):
    """ Sauvegarde les données dans un fichier HDFS """
    try:
        # Vérifier si le répertoire existe, sinon le créer
        if not client.status(HDFS_PATH, strict=False):
            client.makedirs(HDFS_PATH)

        # Écrire les données en mode append
        with client.write(HDFS_FILE_PATH, encoding='utf-8', append=True) as writer:
            writer.write(data + "\n")
        print("Données enregistrées dans HDFS")
    except Exception as e:
        print(f"Erreur lors de l'enregistrement dans HDFS : {e}")

def consumer_kafka():
    """ Lit les messages Kafka et les stocke dans HDFS """
    print("En attente des données Kafka...")
    for message in consumer:
        data = message.value.decode("utf-8")
        print(f"données correctement récuperer: {data}")
        sauvegarder_hdfs(data)

if __name__ == "__main__":
    consumer_kafka()

我有点疑惑,明明配置的HDFS_URL是http://localhost:9870(Namenode的地址),但报错里却指向了hadoop-worker2这个节点,而且我的机器解析不了这个域名。有没有朋友能帮我分析下问题出在哪,该怎么解决呢?

备注:内容来源于stack exchange,提问作者Learning.jUnkie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 19:48:12