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
相关产品推荐
相关产品推荐

