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

Docker环境下Kafka消费数据写入HDFS失败求助

Kafka消费数据无法写入HDFS问题

我正在搭建一套基于Kafka采集数据并存储到HDFS的解决方案,Kafka生产者运行正常,但消费者脚本虽然能消费数据,却无法将数据写入HDFS。

我尝试了两种方法均未成功:

方法一:直接从Kafka Topic写入HDFS

使用以下代码,但执行后无任何反应:

# Define the Kafka consumer
consumer = KafkaConsumer(
    'testing-topic',
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='earliest',
    value_deserializer=lambda x: x.decode('utf-8')
)

# HDFS client 
hdfs_client = InsecureClient('http://localhost:50070', user='root')  

# HDFS file path
hdfs_file_path = '/data/voix_data.csv'

print("Connecting to HDFS...")

try:
    # Open HDFS file for writing
    with hdfs_client.write(hdfs_file_path, encoding='utf-8') as writer:
        # Write CSV header
        writer.write('final_anon_mdn,profile_id,subprofile,product_name,gamme,marche,segment,billing_type,'
                     'final_anon_party_account,final_anon_party_customer,final_anon_contract_id,event_type,'
                     'call_dir,termination_type,network_type,dest_type,other_operator,zone,country,ci,district,'
                     'city,region,num_events,sum_event_minutes,traffic_days,last_timestamp,last_date,start_date,'
                     'end_date,source,id_month\n')
        print("Header written to HDFS.")
        
        for message in consumer:
            # Print the received message
            print(f"Received message: {message.value}")  
            # Write each line to HDFS
            writer.write(message.value + '\n')

方法二:先存储到本地文件再复制到HDFS

执行后报错:copyFromLocal: /data/voix_data.csv': No such file or directory`,使用的脚本如下:

def kafka_to_local_file(kafka_container_id, topic_name, local_file_path):
    # Define the Kafka console consumer command
    kafka_cmd = [
        "docker", "exec", "-i", kafka_container_id,
        "kafka-console-consumer.sh", "--bootstrap-server", "localhost:9092",
        "--topic", topic_name, "--from-beginning", "--timeout-ms", "30000"
    ]
    
    # Open the local file for writing
    with open(local_file_path, 'w', newline='', encoding='utf-8') as local_file:
        csv_writer = csv.writer(local_file)
        
        # Run the Kafka command and capture the output
        process = subprocess.Popen(kafka_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
        
        for line in process.stdout:
            # Decode line and split CSV-like string into list
            decoded_line = line.decode('utf-8').strip()
            print(f"Received line: {decoded_line}")  # Debug print
            if decoded_line:  # Ensure the line is not empty
                row = decoded_line.split(';')
                csv_writer.writerow(row)
        
        process.wait()
        stderr = process.communicate()[1]
        if stderr:
            print(f"Error: {stderr.decode('utf-8')}")

def upload_to_hdfs(local_file_path, hdfs_file_path, hdfs_container_id):
    # Define the HDFS put command
    hdfs_cmd = [
        "docker", "exec", "-i", hdfs_container_id,
        "hdfs", "dfs", "-copyFromLocal", local_file_path, hdfs_file_path
    ]
    
    # Run the HDFS command to upload the file
    result = subprocess.run(hdfs_cmd, capture_output=True, text=True)
    if result.returncode != 0:
        print(f"Failed to upload to HDFS: {result.stderr}")
    else:
        print(f"Data uploaded to HDFS: {result.stdout}")

# Define the Kafka consumer
consumer = KafkaConsumer(
    'testing-topic',  # Topic name
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='earliest',
    value_deserializer=lambda x: x.decode('utf-8')  # Deserialize data from UTF-8 encoded string
)

# HDFS client
hdfs_client = InsecureClient('http://localhost:50070', user='root')

# Local file path and HDFS file path
local_file_path = '/mnt/c/Users/Asus/Desktop/internship/infra/local_test.csv'
hdfs_file_path = '/user/hadoop/data/voix_data.csv'

# Run the functions to transfer data
kafka_container_id = 'bb20bf20183711348208d7a1bf9ae28b69a876419ef79cb1c3f3c66bf5c4d036'
hdfs_container_id = '7bbff9c48e913340cbc97abaa27fc23a8362b2670e714010d75b06458209f4f5'
topic_name = 'testing-topic'  # Replace with your Kafka topic name

print("Starting data transfer from Kafka to local CSV file...")
kafka_to_local_file(kafka_container_id, topic_name, local_file_path)
print("Local CSV file created.")

print("Uploading CSV file to HDFS...")
upload_to_hdfs(local_file_path, hdfs_file_path, hdfs_container_id)

内容的提问来源于stack exchange,提问作者Assou Iben jellal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 09:58:15