Hadoop无法处理中文问题解决方案

2019-03-28 13:08|来源: 网络

由于Hadoop默认编码为UTF-8,并且将UTF-8进行了硬编码,所以我们在处理中文时需要重写OutputFormat类。方法为:

1、新建类GBKFileOutputFormat,代码如下:
import java.io.DataOutputStream; 
import java.io.IOException; 
import java.io.UnsupportedEncodingException; 
 
import org.apache.hadoop.conf.Configuration; 
import org.apache.hadoop.fs.FileSystem; 
import org.apache.hadoop.fs.Path; 
import org.apache.hadoop.fs.FSDataOutputStream;
import org.apache.hadoop.mapreduce.lib.*;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
 
import org.apache.hadoop.io.NullWritable; 
import org.apache.hadoop.io.Text; 
import org.apache.hadoop.io.compress.CompressionCodec; 
import org.apache.hadoop.io.compress.GzipCodec; 
import org.apache.hadoop.mapreduce.OutputFormat; 
import org.apache.hadoop.mapreduce.RecordWriter; 
import org.apache.hadoop.mapreduce.TaskAttemptContext; 
import org.apache.hadoop.util.*; 
 
/** An {@link OutputFormat} that writes plain text files. */ 
public class GBKFileOutputFormat<K, V> extends FileOutputFormat<K, V> {//TextInputFormat是默认的输出文件格式 
  protected static class LineRecordWriter<K, V>//默认 
    extends RecordWriter<K, V> { 
    private static final String utf8 = "GBK";  //硬编码,将“UTF-8”改为“GBK” 
    private static final byte[] newline;//行结束符? 
    static { 
      try { 
        newline = "\n".getBytes(utf8); 
      } catch (UnsupportedEncodingException uee) { 
        throw new IllegalArgumentException("can't find " + utf8 + " encoding"); 
      } 
    } 
 
    protected DataOutputStream out; 
    private final byte[] keyValueSeparator;//key和value的分隔符,默认的好像是Tab 
 
    public LineRecordWriter(DataOutputStream out, String keyValueSeparator) {//构造函数,初始化输出流及分隔符 
      this.out = out; 
      try { 
        this.keyValueSeparator = keyValueSeparator.getBytes(utf8); 
      } catch (UnsupportedEncodingException uee) { 
        throw new IllegalArgumentException("can't find " + utf8 + " encoding"); 
      } 
    } 
 
    public LineRecordWriter(DataOutputStream out) {//默认的分隔符 
      this(out, "\t"); 
    } 
 
    /**
    * Write the object to the byte stream, handling Text as a special输出流是byte格式的
    * case.
    * @param o the object to print是要输出的对象
    * @throws IOException if the write throws, we pass it on
    */ 
    private void writeObject(Object o) throws IOException {//应该是一行一行的写 key keyValueSeparator value \n 
      if (o instanceof Text) {//如果o是Text的实例 
        Text to = (Text) o; 
        out.write(to.getBytes(), 0, to.getLength());//写出 
      } else { 
        out.write(o.toString().getBytes(utf8)); 
      } 
    } 
 
    public synchronized void write(K key, V value)//给写线程加锁,写是互斥行为 
      throws IOException { 
//下面是为了判断key和value是否为空值 
      boolean nullKey = key == null || key instanceof NullWritable;//这语句太牛了 
      boolean nullValue = value == null || value instanceof NullWritable; 
      if (nullKey && nullValue) {// 
        return; 
      } 
      if (!nullKey) { 
        writeObject(key); 
      } 
      if (!(nullKey || nullValue)) { 
        out.write(keyValueSeparator); 
      } 
      if (!nullValue) { 
        writeObject(value); 
      } 
      out.write(newline); 
    } 
 
    public synchronized 
    void close(TaskAttemptContext context) throws IOException { 
      out.close(); 
    } 
  } 
 
  public RecordWriter<K, V>    getRecordWriter(TaskAttemptContext job//获得writer实例 
                        ) throws IOException, InterruptedException { 
    Configuration conf = job.getConfiguration(); 
    boolean isCompressed = getCompressOutput(job);// 
    String keyValueSeparator= conf.get("mapred.textoutputformat.separator", 
                                      "\t"); 
    CompressionCodec codec = null;//压缩格式 还是? 
    String extension = ""; 
    if (isCompressed) { 
      Class<? extends CompressionCodec> codecClass = 
        getOutputCompressorClass(job, GzipCodec.class); 
      codec = (CompressionCodec) ReflectionUtils.newInstance(codecClass, conf); 
      extension = codec.getDefaultExtension(); 
    } 
    Path file = getDefaultWorkFile(job, extension);//这个是获取缺省的文件路径及名称,在FileOutput中有对其的实现 
    FileSystem fs = file.getFileSystem(conf); 
    if (!isCompressed) { 
      FSDataOutputStream fileOut = fs.create(file, false); 
      return new LineRecordWriter<K, V>(fileOut, keyValueSeparator); 
    } else { 
      FSDataOutputStream fileOut = fs.create(file, false); 
      return new LineRecordWriter<K, V>(new DataOutputStream 
                                        (codec.createOutputStream(fileOut)), 
                                        keyValueSeparator); 
    } 
  } 

该类是在源代码中TextOutputFormat类基础上进行修改的,在这需要注意的一点是继承的父类FileOutputFormat是位于org.apache.hadoop.mapreduce.lib.output包中的

2、在主类中添加job.setOutputFormatClass(GBKFileOutputFormat.class);

更多Hadoop相关信息见Hadoop 专题页面 http://www.linuxidc.com/topicnews.aspx?tid=13

相关问答

更多
  • 要上好一节欣赏课,先决条件是教师自身对于名作的深刻解读能力,这就要求教师自身必须具备深厚的艺术修养,那才能真正在课堂上游刃有余。其次江南徐老师的上课风格也给我留下了深刻的映象,在全省的公开课上,教师能如此轻松自然地与学生进行互动探究,感叹教师诙谐幽默的教学风格和自身丰 厚的文化底蕴。相比自己的公开课,就缺少了这样一种状态,但我想这种状态的落实,恰巧需要教师在平常的教学中不断历练积聚。
  • hadoop-examples-1.0.2.jar文件不在当前目录,建议使用全路径表述。比如: hadoop jar /home/hadoop/..../..../hadoop-examples-1.0.2.jar teragen ... ... hadoop jar /home/hadoop/..../..../hadoop-examples-1.0.2.jar terasort ... ... hadoop jar /home/hadoop/..../..../hadoop-examples-1.0.2 ...
  • 用脚本在MySQL更新数据是,传递了不合法的变量。请检查传递的变量是否都有值。 可以在bind_param的时候先判断一下要传递的变量是否defined,如果不是,则传一个undef代表MySQL中的NULL 提交表单的页面里面表单项的ID可能写错了。就是series_id写错了。可以检查一下。
  • 查询语句 $sql="select count(s) as c from stat_20101105 group by s"; show index from stat_20101105结果如下: table->stat_20101105 non_union->0 key_name->primary seq_in_index->1 colunm_name->id callation->a cardinality->7146648 sub_part-> packed-> null-> index_type-> ...
  • 乱码是因为使用的字符集不对。试试修改一下i18n这个文件: vim /etc/sysconfig/i18n,然后按i,将引号里面的内容改为 zh_CN.UTF-8或zh_CN.GBK,注销重登录。
  • 方框跟乱码不一样的,你要记住啊,方框是字体的原因导致的,修改fcitx的配置文件把里面的字体改成系统有的字体。。乱码是因为编码识别错误导致的,不是方块那种。。终端里面的话,你在终端上面的菜单里调一下编码就行了吧,点右键也能调编码的貌似。。
  • 那么,为什么不添加这个标题呢? 只需一行代码:) 您可以运行自己的服务器,将请求转发给该服务(nginx可以轻松地执行此操作, 请查看 )。 简而言之: frontend app -> your proxy with CORS headers (like nginx) -> api service So, why just not add this header? Just one line of code :) You may run your own server which will forward ...
  • 你没有说什么不行,但我会冒险猜测。 你会得到一个空指针例外的地方。 您的代码存在的一个问题是,当您执行onConfigurationChanged时,您需要重复执行onCreate()大部分逻辑。 活动中有一套全新的视图。 否则,成员字段tv将为空,并且不会有听众附加到您的按钮。 编辑根据您的评论,该count没有被保留,我相信您的清单有问题,系统正在销毁并重新创建您的定位更改活动。 要保留这些更改的count ,请覆盖onRetainNonConfigurationInstance()以返回一个包含cou ...
  • 我在这里使用Python切片语法因为它很好。 S[:-1]表示删除了最后一个字符的字符串S , S[:n]表示长度为n的S的前缀。 关键思想是如果C是A和B的交织,则C[:-1]是A和B[:-1]或A[:-1]和B的交织。 另一方面,如果C是A和B的交织,则C + 'X'是A + 'X'和B的交织,它也是A和B + 'X'的交织。 这些是我们应用动态编程所需的子结构属性。 我们定义f(i, j) = true ,如果s1[:i]和s2[:j]可以交织形成s3[:(i+j)] ,否则f(i,j) = fals ...