最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

java kafka寫入數(shù)據(jù)到HDFS問題

 更新時間:2023年08月31日 08:56:35   作者:我是女孩  
這篇文章主要介紹了java kafka寫入數(shù)據(jù)到HDFS問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教

java kafka寫入數(shù)據(jù)到HDFS

安裝kafka,見我以前的文章

http://www.fzitv.net/server/2968144y7.htm

向Hdfs寫入文件,控制臺會輸出以下錯誤信息:

Caused by: org.apache.hadoop.ipc.RemoteException(org.apache.hadoop.security.AccessControlException): Permission denied: user=s00356746, access=WRITE, inode="/user":root:supergroup:drwxr-xr-x

從中很容易看出是因為當前執(zhí)行Spark Application的用戶沒有Hdfs“/user”目錄的寫入權(quán)限。這個問題無論是在Windows下還是Linux下提交Spark Application都經(jīng)常會遇到

如果是歐拉操作系統(tǒng)

需做如下處理

chattr -i etc/passwd
chattr -i /etc/shadow
chattr -i /etc/group
chattr -i /etc/passwd-
chattr -i /etc/shadow-
chattr -i /etc/group-
lsattr passwd*
都需要沒有   i   屬性

如果是Linux環(huán)境

將執(zhí)行操作的用戶添加到supergroup用戶組。

groupadd supergroup
usermod -a -G supergroup s00356746

如果是Windows用戶

在hdfs namenode所在機器添加新用戶,用戶名為執(zhí)行操作的Windows用戶名,然后將此用戶添加到supergroup用戶組。

adduser s00356746
groupadd supergroup
usermod -a -G supergroup s00356746

這樣,以后每次執(zhí)行類似操作可以將文件寫入Hdfs中屬于s00356746用戶的目錄內(nèi),而不會出現(xiàn)上面的Exception。

生產(chǎn)者代碼

import java.util.Properties;
import kafka.javaapi.producer.Producer;
import kafka.producer.KeyedMessage;
import kafka.producer.ProducerConfig;
public class KafkaProducer  {
    private final Producer<String, String> producer;
    public final static String TOPIC = "test";
    private KafkaProducer(){
        Properties props = new Properties();
        //此處配置的是kafka的端口
        props.put("metadata.broker.list", "10.175.118.105:9092");
        //配置value的序列化類
        props.put("serializer.class", "kafka.serializer.StringEncoder");
        //配置key的序列化類
        props.put("key.serializer.class", "kafka.serializer.StringEncoder");
        props.put("request.required.acks","-1");
        producer = new Producer<String, String>(new ProducerConfig(props));
    }
    void produce() {
        int messageNo = 1000;
        final int COUNT = 10000;
        while (messageNo < COUNT) {
            String key = String.valueOf(messageNo);
            String data = "hello kafka message " + key;
            producer.send(new KeyedMessage<String, String>(TOPIC, key ,data));
            System.out.println(data);
            messageNo ++;
        }
    }
    public static void main( String[] args )
    {
        new KafkaProducer().produce();
    }
}

kafka寫入Hdfs

package com.huawei.hwclouds.dbs.ops.huatuo.diagnosis.service.impl;
import com.huawei.hwclouds.dbs.common.exception.DBSErrorCode;
import com.huawei.hwclouds.dbs.common.exception.DBSException;
import com.huawei.hwclouds.dbs.constants.VolumeIoType;
import com.huawei.hwclouds.dbs.coremodel.model.dto.DBSInstanceDto;
import com.huawei.hwclouds.dbs.coremodel.model.dto.DBSNodeDto;
import com.huawei.hwclouds.dbs.coremodel.resource.dto.DBSResourceSpecDto;
import com.huawei.hwclouds.dbs.coremodel.resource.dto.DBSVolumeDto;
import kafka.consumer.Consumer;
import kafka.consumer.ConsumerConfig;
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.javaapi.consumer.ConsumerConnector;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataOutputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IOUtils;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.text.SimpleDateFormat;
import java.util.Calendar;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.stream.Collectors;
public class KafkaToHdfs extends Thread {
    private static String kafkaHost = null;
    private static String kafkaGroup = null;
    private static String kafkaTopic = null;
    private static String hdfsUri = null;
    private static String hdfsDir = null;
    private static String hadoopUser = null;
    private static Boolean isDebug = false;
    private ConsumerConnector consumer = null;
    private static Configuration hdfsConf = null;
    private static FileSystem hadoopFS = null;
    public static void main(String[] args) {
//        if (args.length < 6) {
//            useage();
//            System.exit(0);
//        }
//        Map<String, String> user = new HashMap<String, String>();
//        user = System.getenv();
//        user.put("HADOOP_USER_NAME","hadoop");
//        if (user.get("HADOOP_USER_NAME") == null) {
//            System.out.println("請設(shè)定hadoop的啟動的用戶名,環(huán)境變量名稱:HADOOP_USER_NAME,對應(yīng)的值是hadoop的啟動的用戶名");
//            System.exit(0);
//        } else {
//            hadoopUser = user.get("HADOOP_USER_NAME");
//        }
        hadoopUser = "hadoop";
        init(args);
        System.out.println("開始啟動服務(wù)...");
        hdfsConf = new Configuration();
        try {
            hdfsConf.set("fs.defaultFS", hdfsUri);
            hdfsConf.set("dfs.support.append", "true");
            hdfsConf.set("dfs.client.block.write.replace-datanode-on-failure.policy", "NEVER");
            hdfsConf.set("dfs.client.block.write.replace-datanode-on-failure.enable", "true");
        } catch (Exception e) {
            System.out.println(e);
        }
        //創(chuàng)建好相應(yīng)的目錄
        try {
            hadoopFS = FileSystem.get(hdfsConf);
            //如果hdfs的對應(yīng)的目錄不存在,則進行創(chuàng)建
            if (!hadoopFS.exists(new Path("/" + hdfsDir))) {
                hadoopFS.mkdirs(new Path("/" + hdfsDir));
            }
            hadoopFS.close();
        } catch (IOException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
        KafkaToHdfs selfObj = new KafkaToHdfs();
        selfObj.start();
        System.out.println("服務(wù)啟動完畢,監(jiān)聽執(zhí)行中");
    }
    public void run() {
        Properties props = new Properties();
        props.put("zookeeper.connect", kafkaHost);
        props.put("group.id", kafkaGroup);
        props.put("zookeeper.session.timeout.ms", "10000");
        props.put("zookeeper.sync.time.ms", "200");
        props.put("auto.commit.interval.ms", "1000");
        props.put("auto.offset.reset", "smallest");
        props.put("format", "binary");
        props.put("auto.commit.enable", "true");
        props.put("serializer.class", "kafka.serializer.StringEncoder");
        ConsumerConfig consumerConfig = new ConsumerConfig(props);
        this.consumer = Consumer.createJavaConsumerConnector(consumerConfig);
        Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
        topicCountMap.put(kafkaTopic, new Integer(1));
        Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap);
        KafkaStream<byte[], byte[]> stream = consumerMap.get(kafkaTopic).get(0);
        ConsumerIterator<byte[], byte[]> it = stream.iterator();
        while (it.hasNext()) {
            String tmp = new String(it.next().message());
            String fileContent = null;
            if (!tmp.endsWith("\n"))
                fileContent = new String(tmp + "\n");
            else
                fileContent = tmp;
            debug("receive: " + fileContent);
            try {
                hadoopFS = FileSystem.get(hdfsConf);
                String fileName = "/" + hdfsDir + "/" +
                        (new SimpleDateFormat("yyyy-MM-dd").format(Calendar.getInstance().getTime())) + ".txt";
                Path dst = new Path(fileName);
                if (!hadoopFS.exists(dst)) {
                    FSDataOutputStream output = hadoopFS.create(dst);
                    output.close();
                }
                InputStream in = new ByteArrayInputStream(fileContent.getBytes("UTF-8"));
                OutputStream out = hadoopFS.append(dst);
                IOUtils.copyBytes(in, out, 4096, true);
            } catch (IOException e) {
                e.printStackTrace();
            } finally {
                try {
                    hadoopFS.close();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }
        }
        consumer.shutdown();
    }
    private static void init(String[] args) {
        kafkaHost = "10.175.118.105:2182";
        kafkaGroup = "test-consumer-group";
        kafkaTopic = "test";
        hdfsUri = "hdfs://10.175.118.105:9000";
        hdfsDir = "shxsh";
        if (args.length > 5) {
            if (args[5].equals("true")) {
                isDebug = true;
            }
        }
        debug("初始化服務(wù)參數(shù)完畢,參數(shù)信息如下");
        debug("KAFKA_HOST: " + kafkaHost);
        debug("KAFKA_GROUP: " + kafkaGroup);
        debug("KAFKA_TOPIC: " + kafkaTopic);
        debug("HDFS_URI: " + hdfsUri);
        debug("HDFS_DIRECTORY: " + hdfsDir);
        debug("HADOOP_USER: " + hadoopUser);
        debug("IS_DEBUG: " + isDebug);
    }
    private static void debug(String str) {
        if (isDebug) {
            System.out.println(str);
        }
    }
    private static void useage() {
        System.out.println("* kafka寫入到hdfs的Java工具使用說明 ");
        System.out.println("# java -cp kafkatohdfs.jar KafkaToHdfs KAFKA_HOST KAFKA_GROUP KAFKA_TOPIC HDFS_URI HDFS_DIRECTORY IS_DEBUG");
        System.out.println("*  參數(shù)說明:");
        System.out.println("*   KAFKA_HOST      : 代表kafka的主機名或IP:port,例如:namenode:2181,datanode1:2181,datanode2:2181");
        System.out.println("*   KAFKA_GROUP     : 代表kafka的組,例如:test-consumer-group");
        System.out.println("*   KAFKA_TOPIC     : 代表kafka的topic名稱 ,例如:usertags");
        System.out.println("*   HDFS_URI        : 代表hdfs鏈接uri ,例如:hdfs://namenode:9000");
        System.out.println("*   HDFS_DIRECTORY  : 代表hdfs目錄名稱 ,例如:usertags");
        System.out.println("*  可選參數(shù):");
        System.out.println("*   IS_DEBUG        : 代表是否開啟調(diào)試模式,true是,false否,默認為false");
    }
}
 

總結(jié)

以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Gradle中Maven倉庫配置的實現(xiàn)步驟

    Gradle中Maven倉庫配置的實現(xiàn)步驟

    在Java/Kotlin 項目的構(gòu)建過程中,依賴管理是核心環(huán)節(jié)之一,而Maven 倉庫是最常見的依賴來源,本文就來詳細的介紹一下Gradle中Maven倉庫配置,感興趣的可以了解一下
    2025-09-09
  • 微信公眾平臺(測試接口)準備工作

    微信公眾平臺(測試接口)準備工作

    想要微信開發(fā),首先要有個服務(wù)器,但是自己沒有。這時候可以用花生殼,將內(nèi)網(wǎng)映射到公網(wǎng)上,這樣就可以在公網(wǎng)訪問自己的網(wǎng)站了。
    2016-05-05
  • mybatis動態(tài)插入list傳入List參數(shù)的實例代碼

    mybatis動態(tài)插入list傳入List參數(shù)的實例代碼

    本文通過實例代碼給大家介紹了mybatis動態(tài)插入list,Mybatis 傳入List參數(shù)的方法,非常不錯,具有參考借鑒價值,需要的朋友參考下吧
    2018-04-04
  • SpringBoot利用攔截器實現(xiàn)避免重復(fù)請求

    SpringBoot利用攔截器實現(xiàn)避免重復(fù)請求

    Spring MVC中的攔截器(Interceptor)類似于Servlet中的過濾器(Filter),它主要用于攔截用戶請求并作相應(yīng)的處理。本文就將利用攔截器實現(xiàn)避免重復(fù)請求,感興趣的小伙伴可以了解一下
    2022-11-11
  • spring初始化源碼代碼淺析

    spring初始化源碼代碼淺析

    Spring框架被廣泛應(yīng)用于我們的日常工作中,但是很長時間以來我們都是只會使用,不懂它的作用原理,下面這篇文章主要給大家介紹了關(guān)于spring初始化源碼的相關(guān)資料,需要的朋友可以參考下
    2023-04-04
  • Java實現(xiàn)線程按序交替執(zhí)行的方法詳解

    Java實現(xiàn)線程按序交替執(zhí)行的方法詳解

    這篇文章主要為大家詳細介紹了Java如何實現(xiàn)線程按序交替執(zhí)行,文中的示例代碼講解詳細,對我們了解線程有一定幫助,需要的可以參考一下
    2022-10-10
  • eclipse自動提示和自動補全功能實現(xiàn)方法

    eclipse自動提示和自動補全功能實現(xiàn)方法

    這篇文章主要介紹了eclipse自動提示和自動補全的相關(guān)內(nèi)容,文中向大家分享了二者的實現(xiàn)方法代碼,需要的朋友可以了解下。
    2017-09-09
  • Spark操作之a(chǎn)ggregate、aggregateByKey詳解

    Spark操作之a(chǎn)ggregate、aggregateByKey詳解

    這篇文章主要介紹了Spark操作之a(chǎn)ggregate、aggregateByKey詳解,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • IntelliJ IDEA如何設(shè)置JDK版本

    IntelliJ IDEA如何設(shè)置JDK版本

    這篇文章主要介紹了IntelliJ IDEA如何設(shè)置JDK版本問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-12-12
  • Java OpenCV圖像處理之圖形與文字繪制

    Java OpenCV圖像處理之圖形與文字繪制

    這篇文章主要為大家介紹了如何利益Java OpenCV實現(xiàn)在圖像上繪制文字與形狀,文中的示例代碼講解詳細,感興趣的小伙伴可以跟隨小編一起動手試一試
    2022-02-02

最新評論

南开区| 韶关市| 浦北县| 北安市| 曲靖市| 天津市| 仙游县| 达孜县| 壤塘县| 安义县| 淅川县| 广灵县| 张家川| 靖西县| 南阳市| 汽车| 隆昌县| 浦县| 中西区| 阳曲县| 青岛市| 东方市| 苏尼特左旗| 金沙县| 定州市| 南安市| 邹城市| 中西区| 睢宁县| 神农架林区| 卢氏县| 苏尼特左旗| 遂溪县| 湖南省| 唐河县| 汝州市| 广丰县| 连云港市| 体育| 南皮县| 专栏|