`
weitao1026
  • 浏览: 1055932 次
  • 性别: Icon_minigender_1
  • 来自: 上海
社区版块
存档分类
最新评论
阅读更多

如何使用Apache Pig与Lucene集成,还不知道的道友们,可以先看下上篇,熟悉下具体的流程。
在与Lucene集成过程中,我们发现最终还要把生成的Lucene索引,拷贝至本地磁盘,才能提供检索服务,这样以来,比较繁琐,而且有以下几个缺点:

(一)在生成索引以及最终能提供正常的服务之前,索引经过多次落地操作,这无疑会给磁盘和网络IO,带来巨大影响

(二)Lucene的Field的配置与其UDF函数的代码耦合性过强,而且提供的配置也比较简单,不太容易满足,灵活多变的检索需求和服务,如果改动索引配置,则有可能需要重新编译源码。

(三)对Hadoop的分布式存储系统HDFS依赖过强,如果使用与Lucene集成,那么则意味着你提供检索的Web服务器,则必须跟hadoop的存储节点在一个机器上,否则,无法从HDFS上下拉索引,除非你自己写程序,或使用scp再次从目标机传输,这样无疑又增加了,系统的复杂性。


鉴于有以上几个缺点,所以建议大家使用Solr或ElasticSearch这样的封装了Lucene更高级的API框架,那么Solr与ElasticSearch和Lucene相比,又有什么优点呢?

(1)在最终的写入数据时,我们可以直接最终结果写入solr或es,同时也可以在HDFS上保存一份,作为灾备。

(2)使用了solr或es,这时,我们字段的配置完全与UDF函数代码无关,我们的任何字段配置的变动,都不会影响Pig的UDF函数的代码,而在UDF函数里,唯一要做的,就是将最终数据,提供给solr和es服务。

(3)solr和es都提供了restful风格的http操作方式,这时候,我们的检索集群完全可以与Hadoop集群分离,从而让他们各自都专注自己的服务。



下面,散仙就具体说下如何使用Pig和Solr集成?

(1)依旧访问这个地址下载源码压缩包。
(2)提取出自己想要的部分,在eclipse工程中,修改定制适合自己环境的的代码(Solr版本是否兼容?hadoop版本是否兼容?,Pig版本是否兼容?)。
(3)使用ant重新打包成jar
(4)在pig里,注册相关依赖的jar包,并使用索引存储



注意,在github下载的压缩里直接提供了对SolrCloud模式的提供,而没有提供,普通模式的函数,散仙在这里稍作修改后,可以支持普通模式的Solr服务,代码如下:


SolrOutputFormat函数

Java代码 复制代码 收藏代码
  1. package com.pig.support.solr;  
  2.   
  3.   
  4.   
  5. import java.io.IOException;  
  6. import java.util.ArrayList;  
  7. import java.util.List;  
  8. import java.util.concurrent.Executors;  
  9. import java.util.concurrent.ScheduledExecutorService;  
  10. import java.util.concurrent.TimeUnit;  
  11.   
  12. import org.apache.hadoop.io.Writable;  
  13. import org.apache.hadoop.mapreduce.JobContext;  
  14. import org.apache.hadoop.mapreduce.OutputCommitter;  
  15. import org.apache.hadoop.mapreduce.RecordWriter;  
  16. import org.apache.hadoop.mapreduce.TaskAttemptContext;  
  17. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;  
  18. import org.apache.solr.client.solrj.SolrServer;  
  19. import org.apache.solr.client.solrj.SolrServerException;  
  20. import org.apache.solr.client.solrj.impl.CloudSolrServer;  
  21. import org.apache.solr.client.solrj.impl.HttpSolrServer;  
  22. import org.apache.solr.common.SolrInputDocument;  
  23. /** 
  24.  * @author qindongliang 
  25.  * 支持SOlr的SolrOutputFormat 
  26.  * 如果你想了解,或学习更多这方面的 
  27.  * 知识,请加入我们的群: 
  28.  *  
  29.  * 搜索技术交流群(2000人):324714439  
  30.  * 大数据技术1号交流群(2000人):376932160  (已满) 
  31.  * 大数据技术2号交流群(2000人):415886155  
  32.  * 微信公众号:我是攻城师(woshigcs) 
  33.  *  
  34.  * */  
  35. public class SolrOutputFormat extends  
  36.         FileOutputFormat<Writable, SolrInputDocument> {  
  37.   
  38.     final String address;  
  39.     final String collection;  
  40.   
  41.     public SolrOutputFormat(String address, String collection) {  
  42.         this.address = address;  
  43.         this.collection = collection;  
  44.     }  
  45.   
  46.     @Override  
  47.     public RecordWriter<Writable, SolrInputDocument> getRecordWriter(  
  48.             TaskAttemptContext ctx) throws IOException, InterruptedException {  
  49.         return new SolrRecordWriter(ctx, address, collection);  
  50.     }  
  51.   
  52.       
  53.     @Override  
  54.     public synchronized OutputCommitter getOutputCommitter(  
  55.             TaskAttemptContext arg0) throws IOException {  
  56.         return new OutputCommitter(){  
  57.   
  58.             @Override  
  59.             public void abortTask(TaskAttemptContext ctx) throws IOException {  
  60.                   
  61.             }  
  62.   
  63.             @Override  
  64.             public void commitTask(TaskAttemptContext ctx) throws IOException {  
  65.                   
  66.             }  
  67.   
  68.             @Override  
  69.             public boolean needsTaskCommit(TaskAttemptContext arg0)  
  70.                     throws IOException {  
  71.                 return true;  
  72.             }  
  73.   
  74.             @Override  
  75.             public void setupJob(JobContext ctx) throws IOException {  
  76.                   
  77.             }  
  78.   
  79.             @Override  
  80.             public void setupTask(TaskAttemptContext ctx) throws IOException {  
  81.                   
  82.             }  
  83.               
  84.               
  85.         };  
  86.     }  
  87.   
  88.   
  89.     /** 
  90.      * Write out the LuceneIndex to a local temporary location.<br/> 
  91.      * On commit/close the index is copied to the hdfs output directory.<br/> 
  92.      *  
  93.      */  
  94.     static class SolrRecordWriter extends RecordWriter<Writable, SolrInputDocument> {  
  95.         /**Solr的地址*/  
  96.         SolrServer server;  
  97.         /**批处理提交的数量**/  
  98.         int batch = 5000;  
  99.           
  100.         TaskAttemptContext ctx;  
  101.           
  102.         List<SolrInputDocument> docs = new ArrayList<SolrInputDocument>(batch);  
  103.         ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor();  
  104.         /** 
  105.          * Opens and forces connect to CloudSolrServer 
  106.          *  
  107.          * @param address 
  108.          */  
  109.         public SolrRecordWriter(final TaskAttemptContext ctx, String address, String collection) {  
  110.             try {  
  111.                 this.ctx = ctx;  
  112.                 server = new HttpSolrServer(address);  
  113.                   
  114.                 exec.scheduleWithFixedDelay(new Runnable(){  
  115.                     public void run(){  
  116.                         ctx.progress();  
  117.                     }  
  118.                 }, 10001000, TimeUnit.MILLISECONDS);  
  119.             } catch (Exception e) {  
  120.                 RuntimeException exc = new RuntimeException(e.toString(), e);  
  121.                 exc.setStackTrace(e.getStackTrace());  
  122.                 throw exc;  
  123.             }  
  124.         }  
  125.   
  126.           
  127.         /** 
  128.          * On close we commit 
  129.          */  
  130.         @Override  
  131.         public void close(final TaskAttemptContext ctx) throws IOException,  
  132.                 InterruptedException {  
  133.   
  134.             try {  
  135.                   
  136.                 if (docs.size() > 0) {  
  137.                     server.add(docs);  
  138.                     docs.clear();  
  139.                 }  
  140.   
  141.                 server.commit();  
  142.             } catch (SolrServerException e) {  
  143.                 RuntimeException exc = new RuntimeException(e.toString(), e);  
  144.                 exc.setStackTrace(e.getStackTrace());  
  145.                 throw exc;  
  146.             } finally {  
  147.                 server.shutdown();  
  148.                 exec.shutdownNow();  
  149.             }  
  150.               
  151.         }  
  152.   
  153.         /** 
  154.          * We add the indexed documents without commit 
  155.          */  
  156.         @Override  
  157.         public void write(Writable key, SolrInputDocument doc)  
  158.                 throws IOException, InterruptedException {  
  159.             try {  
  160.                 docs.add(doc);  
  161.                 if (docs.size() >= batch) {  
  162.                     server.add(docs);  
  163.                     docs.clear();  
  164.                 }  
  165.             } catch (SolrServerException e) {  
  166.                 RuntimeException exc = new RuntimeException(e.toString(), e);  
  167.                 exc.setStackTrace(e.getStackTrace());  
  168.                 throw exc;  
  169.             }  
  170.         }  
  171.   
  172.     }  
  173. }  
package com.pig.support.solr;



import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;

import org.apache.hadoop.io.Writable;
import org.apache.hadoop.mapreduce.JobContext;
import org.apache.hadoop.mapreduce.OutputCommitter;
import org.apache.hadoop.mapreduce.RecordWriter;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.solr.client.solrj.SolrServer;
import org.apache.solr.client.solrj.SolrServerException;
import org.apache.solr.client.solrj.impl.CloudSolrServer;
import org.apache.solr.client.solrj.impl.HttpSolrServer;
import org.apache.solr.common.SolrInputDocument;
/**
 * @author qindongliang
 * 支持SOlr的SolrOutputFormat
 * 如果你想了解,或学习更多这方面的
 * 知识,请加入我们的群:
 * 
 * 搜索技术交流群(2000人):324714439 
 * 大数据技术1号交流群(2000人):376932160  (已满)
 * 大数据技术2号交流群(2000人):415886155 
 * 微信公众号:我是攻城师(woshigcs)
 * 
 * */
public class SolrOutputFormat extends
		FileOutputFormat<Writable, SolrInputDocument> {

	final String address;
	final String collection;

	public SolrOutputFormat(String address, String collection) {
		this.address = address;
		this.collection = collection;
	}

	@Override
	public RecordWriter<Writable, SolrInputDocument> getRecordWriter(
			TaskAttemptContext ctx) throws IOException, InterruptedException {
		return new SolrRecordWriter(ctx, address, collection);
	}

	
	@Override
	public synchronized OutputCommitter getOutputCommitter(
			TaskAttemptContext arg0) throws IOException {
		return new OutputCommitter(){

			@Override
			public void abortTask(TaskAttemptContext ctx) throws IOException {
				
			}

			@Override
			public void commitTask(TaskAttemptContext ctx) throws IOException {
				
			}

			@Override
			public boolean needsTaskCommit(TaskAttemptContext arg0)
					throws IOException {
				return true;
			}

			@Override
			public void setupJob(JobContext ctx) throws IOException {
				
			}

			@Override
			public void setupTask(TaskAttemptContext ctx) throws IOException {
				
			}
			
			
		};
	}


	/**
	 * Write out the LuceneIndex to a local temporary location.<br/>
	 * On commit/close the index is copied to the hdfs output directory.<br/>
	 * 
	 */
	static class SolrRecordWriter extends RecordWriter<Writable, SolrInputDocument> {
		/**Solr的地址*/
		SolrServer server;
		/**批处理提交的数量**/
		int batch = 5000;
		
		TaskAttemptContext ctx;
		
		List<SolrInputDocument> docs = new ArrayList<SolrInputDocument>(batch);
		ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor();
		/**
		 * Opens and forces connect to CloudSolrServer
		 * 
		 * @param address
		 */
		public SolrRecordWriter(final TaskAttemptContext ctx, String address, String collection) {
			try {
				this.ctx = ctx;
				server = new HttpSolrServer(address);
				
				exec.scheduleWithFixedDelay(new Runnable(){
					public void run(){
						ctx.progress();
					}
				}, 1000, 1000, TimeUnit.MILLISECONDS);
			} catch (Exception e) {
				RuntimeException exc = new RuntimeException(e.toString(), e);
				exc.setStackTrace(e.getStackTrace());
				throw exc;
			}
		}

		
		/**
		 * On close we commit
		 */
		@Override
		public void close(final TaskAttemptContext ctx) throws IOException,
				InterruptedException {

			try {
				
				if (docs.size() > 0) {
					server.add(docs);
					docs.clear();
				}

				server.commit();
			} catch (SolrServerException e) {
				RuntimeException exc = new RuntimeException(e.toString(), e);
				exc.setStackTrace(e.getStackTrace());
				throw exc;
			} finally {
				server.shutdown();
				exec.shutdownNow();
			}
			
		}

		/**
		 * We add the indexed documents without commit
		 */
		@Override
		public void write(Writable key, SolrInputDocument doc)
				throws IOException, InterruptedException {
			try {
				docs.add(doc);
				if (docs.size() >= batch) {
					server.add(docs);
					docs.clear();
				}
			} catch (SolrServerException e) {
				RuntimeException exc = new RuntimeException(e.toString(), e);
				exc.setStackTrace(e.getStackTrace());
				throw exc;
			}
		}

	}
}





SolrStore函数

Java代码 复制代码 收藏代码
  1. package com.pig.support.solr;  
  2.   
  3.   
  4.   
  5. import java.io.IOException;  
  6. import java.util.Properties;  
  7.   
  8. import org.apache.hadoop.fs.Path;  
  9. import org.apache.hadoop.io.Writable;  
  10. import org.apache.hadoop.mapreduce.Job;  
  11. import org.apache.hadoop.mapreduce.OutputFormat;  
  12. import org.apache.hadoop.mapreduce.RecordWriter;  
  13. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;  
  14. import org.apache.pig.ResourceSchema;  
  15. import org.apache.pig.ResourceSchema.ResourceFieldSchema;  
  16. import org.apache.pig.ResourceStatistics;  
  17. import org.apache.pig.StoreFunc;  
  18. import org.apache.pig.StoreMetadata;  
  19. import org.apache.pig.data.Tuple;  
  20. import org.apache.pig.impl.util.UDFContext;  
  21. import org.apache.pig.impl.util.Utils;  
  22. import org.apache.solr.common.SolrInputDocument;  
  23.   
  24. /** 
  25.  *  
  26.  * Create a lucene index 
  27.  *  
  28.  */  
  29. public class SolrStore extends StoreFunc implements StoreMetadata {  
  30.   
  31.     private static final String SCHEMA_SIGNATURE = "solr.output.schema";  
  32.   
  33.     ResourceSchema schema;  
  34.     String udfSignature;  
  35.     RecordWriter<Writable, SolrInputDocument> writer;  
  36.   
  37.     String address;  
  38.     String collection;  
  39.       
  40.     public SolrStore(String address, String collection) {  
  41.         this.address = address;  
  42.         this.collection = collection;  
  43.     }  
  44.   
  45.     public void storeStatistics(ResourceStatistics stats, String location,  
  46.             Job job) throws IOException {  
  47.     }  
  48.   
  49.     public void storeSchema(ResourceSchema schema, String location, Job job)  
  50.             throws IOException {  
  51.     }  
  52.   
  53.     @Override  
  54.     public void checkSchema(ResourceSchema s) throws IOException {  
  55.         UDFContext udfc = UDFContext.getUDFContext();  
  56.         Properties p = udfc.getUDFProperties(this.getClass(),  
  57.                 new String[] { udfSignature });  
  58.         p.setProperty(SCHEMA_SIGNATURE, s.toString());  
  59.     }  
  60.   
  61.     public OutputFormat<Writable, SolrInputDocument> getOutputFormat()  
  62.             throws IOException {  
  63.         // not be used  
  64.         return new SolrOutputFormat(address, collection);  
  65.     }  
  66.   
  67.     /** 
  68.      * Not used 
  69.      */  
  70.     @Override  
  71.     public void setStoreLocation(String location, Job job) throws IOException {  
  72.         FileOutputFormat.setOutputPath(job, new Path(location));  
  73.     }  
  74.   
  75.     @Override  
  76.     public void setStoreFuncUDFContextSignature(String signature) {  
  77.         this.udfSignature = signature;  
  78.     }  
  79.   
  80.     @SuppressWarnings({ "unchecked""rawtypes" })  
  81.     @Override  
  82.     public void prepareToWrite(RecordWriter writer) throws IOException {  
  83.         this.writer = writer;  
  84.         UDFContext udc = UDFContext.getUDFContext();  
  85.         String schemaStr = udc.getUDFProperties(this.getClass(),  
  86.                 new String[] { udfSignature }).getProperty(SCHEMA_SIGNATURE);  
  87.   
  88.         if (schemaStr == null) {  
  89.             throw new RuntimeException("Could not find udf signature");  
  90.         }  
  91.   
  92.         schema = new ResourceSchema(Utils.getSchemaFromString(schemaStr));  
  93.   
  94.     }  
  95.   
  96.     /** 
  97.      * Shamelessly copied from : https://issues.apache.org/jira/secure/attachment/12484764/NUTCH-1016-2.0.patch 
  98.      * @param input 
  99.      * @return 
  100.      */  
  101.     private static String stripNonCharCodepoints(String input) {  
  102.         StringBuilder retval = new StringBuilder(input.length());  
  103.         char ch;  
  104.   
  105.         for (int i = 0; i < input.length(); i++) {  
  106.             ch = input.charAt(i);  
  107.   
  108.             // Strip all non-characters  
  109.             // http://unicode.org/cldr/utility/list-unicodeset.jsp?a=[:Noncharacter_Code_Point=True:]  
  110.             // and non-printable control characters except tabulator, new line  
  111.             // and carriage return  
  112.             if (ch % 0x10000 != 0xffff && // 0xffff - 0x10ffff range step  
  113.                                             // 0x10000  
  114.                     ch % 0x10000 != 0xfffe && // 0xfffe - 0x10fffe range  
  115.                     (ch <= 0xfdd0 || ch >= 0xfdef) && // 0xfdd0 - 0xfdef  
  116.                     (ch > 0x1F || ch == 0x9 || ch == 0xa || ch == 0xd)) {  
  117.   
  118.                 retval.append(ch);  
  119.             }  
  120.         }  
  121.   
  122.         return retval.toString();  
  123.     }  
  124.   
  125.     @Override  
  126.     public void putNext(Tuple t) throws IOException {  
  127.   
  128.         final SolrInputDocument doc = new SolrInputDocument();  
  129.   
  130.         final ResourceFieldSchema[] fields = schema.getFields();  
  131.         int docfields = 0;  
  132.   
  133.         for (int i = 0; i < fields.length; i++) {  
  134.             final Object value = t.get(i);  
  135.   
  136.             if (value != null) {  
  137.                 docfields++;  
  138.                 doc.addField(fields[i].getName().trim(), stripNonCharCodepoints(value.toString()));  
  139.             }  
  140.   
  141.         }  
  142.   
  143.         try {  
  144.             if (docfields > 0)  
  145.                 writer.write(null, doc);  
  146.         } catch (InterruptedException e) {  
  147.             Thread.currentThread().interrupt();  
  148.             return;  
  149.         }  
  150.   
  151.     }  
  152.   
  153. }  
package com.pig.support.solr;



import java.io.IOException;
import java.util.Properties;

import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Writable;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.OutputFormat;
import org.apache.hadoop.mapreduce.RecordWriter;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.pig.ResourceSchema;
import org.apache.pig.ResourceSchema.ResourceFieldSchema;
import org.apache.pig.ResourceStatistics;
import org.apache.pig.StoreFunc;
import org.apache.pig.StoreMetadata;
import org.apache.pig.data.Tuple;
import org.apache.pig.impl.util.UDFContext;
import org.apache.pig.impl.util.Utils;
import org.apache.solr.common.SolrInputDocument;

/**
 * 
 * Create a lucene index
 * 
 */
public class SolrStore extends StoreFunc implements StoreMetadata {

	private static final String SCHEMA_SIGNATURE = "solr.output.schema";

	ResourceSchema schema;
	String udfSignature;
	RecordWriter<Writable, SolrInputDocument> writer;

	String address;
	String collection;
	
	public SolrStore(String address, String collection) {
		this.address = address;
		this.collection = collection;
	}

	public void storeStatistics(ResourceStatistics stats, String location,
			Job job) throws IOException {
	}

	public void storeSchema(ResourceSchema schema, String location, Job job)
			throws IOException {
	}

	@Override
	public void checkSchema(ResourceSchema s) throws IOException {
		UDFContext udfc = UDFContext.getUDFContext();
		Properties p = udfc.getUDFProperties(this.getClass(),
				new String[] { udfSignature });
		p.setProperty(SCHEMA_SIGNATURE, s.toString());
	}

	public OutputFormat<Writable, SolrInputDocument> getOutputFormat()
			throws IOException {
		// not be used
		return new SolrOutputFormat(address, collection);
	}

	/**
	 * Not used
	 */
	@Override
	public void setStoreLocation(String location, Job job) throws IOException {
		FileOutputFormat.setOutputPath(job, new Path(location));
	}

	@Override
	public void setStoreFuncUDFContextSignature(String signature) {
		this.udfSignature = signature;
	}

	@SuppressWarnings({ "unchecked", "rawtypes" })
	@Override
	public void prepareToWrite(RecordWriter writer) throws IOException {
		this.writer = writer;
		UDFContext udc = UDFContext.getUDFContext();
		String schemaStr = udc.getUDFProperties(this.getClass(),
				new String[] { udfSignature }).getProperty(SCHEMA_SIGNATURE);

		if (schemaStr == null) {
			throw new RuntimeException("Could not find udf signature");
		}

		schema = new ResourceSchema(Utils.getSchemaFromString(schemaStr));

	}

	/**
	 * Shamelessly copied from : https://issues.apache.org/jira/secure/attachment/12484764/NUTCH-1016-2.0.patch
	 * @param input
	 * @return
	 */
	private static String stripNonCharCodepoints(String input) {
		StringBuilder retval = new StringBuilder(input.length());
		char ch;

		for (int i = 0; i < input.length(); i++) {
			ch = input.charAt(i);

			// Strip all non-characters
			// http://unicode.org/cldr/utility/list-unicodeset.jsp?a=[:Noncharacter_Code_Point=True:]
			// and non-printable control characters except tabulator, new line
			// and carriage return
			if (ch % 0x10000 != 0xffff && // 0xffff - 0x10ffff range step
											// 0x10000
					ch % 0x10000 != 0xfffe && // 0xfffe - 0x10fffe range
					(ch <= 0xfdd0 || ch >= 0xfdef) && // 0xfdd0 - 0xfdef
					(ch > 0x1F || ch == 0x9 || ch == 0xa || ch == 0xd)) {

				retval.append(ch);
			}
		}

		return retval.toString();
	}

	@Override
	public void putNext(Tuple t) throws IOException {

		final SolrInputDocument doc = new SolrInputDocument();

		final ResourceFieldSchema[] fields = schema.getFields();
		int docfields = 0;

		for (int i = 0; i < fields.length; i++) {
			final Object value = t.get(i);

			if (value != null) {
				docfields++;
				doc.addField(fields[i].getName().trim(), stripNonCharCodepoints(value.toString()));
			}

		}

		try {
			if (docfields > 0)
				writer.write(null, doc);
		} catch (InterruptedException e) {
			Thread.currentThread().interrupt();
			return;
		}

	}

}



Pig脚本如下:

Java代码 复制代码 收藏代码
  1. --注册依赖文件的jar包  
  2. REGISTER ./dependfiles/tools.jar;  
  3.   
  4. --注册solr相关的jar包  
  5. REGISTER  ./solrdependfiles/pigudf.jar;   
  6. REGISTER  ./solrdependfiles/solr-core-4.10.2.jar;  
  7. REGISTER  ./solrdependfiles/solr-solrj-4.10.2.jar;  
  8. REGISTER  ./solrdependfiles/httpclient-4.3.1.jar  
  9. REGISTER  ./solrdependfiles/httpcore-4.3.jar  
  10. REGISTER  ./solrdependfiles/httpmime-4.3.1.jar  
  11. REGISTER  ./solrdependfiles/noggit-0.5.jar  
  12.   
  13.   
  14. --加载HDFS数据,并定义scheaml  
  15. a = load '/tmp/data' using PigStorage(',') as (sword:chararray,scount:int);  
  16.   
  17. --存储到solr中,并提供solr的ip地址和端口号  
  18. store d into '/user/search/solrindextemp'  using com.pig.support.solr.SolrStore('http://localhost:8983/solr/collection1','collection1');  
  19. ~                                                                                                                                                              
  20. ~                                                                        
  21. ~                                 
--注册依赖文件的jar包
REGISTER ./dependfiles/tools.jar;

--注册solr相关的jar包
REGISTER  ./solrdependfiles/pigudf.jar; 
REGISTER  ./solrdependfiles/solr-core-4.10.2.jar;
REGISTER  ./solrdependfiles/solr-solrj-4.10.2.jar;
REGISTER  ./solrdependfiles/httpclient-4.3.1.jar
REGISTER  ./solrdependfiles/httpcore-4.3.jar
REGISTER  ./solrdependfiles/httpmime-4.3.1.jar
REGISTER  ./solrdependfiles/noggit-0.5.jar


--加载HDFS数据,并定义scheaml
a = load '/tmp/data' using PigStorage(',') as (sword:chararray,scount:int);

--存储到solr中,并提供solr的ip地址和端口号
store d into '/user/search/solrindextemp'  using com.pig.support.solr.SolrStore('http://localhost:8983/solr/collection1','collection1');
~                                                                                                                                                            
~                                                                      
~                               



配置成功之后,我们就可以运行程序,加载HDFS上数据,经过计算处理之后,并将最终的结果,存储到Solr之中,截图如下:






成功之后,我们就可以很方便的在solr中进行毫秒级别的操作了,例如各种各样的全文查询,过滤,排序统计等等!

同样的方式,我们也可以将索引存储在ElasticSearch中,关于如何使用Pig和ElasticSearch集成,散仙也会在后面的文章中介绍,敬请期待!

分享到:
评论

相关推荐

    elasticsearch与hadoop比较

    而Hadoop是一个由Apache软件基金会开发的开源框架,它允许使用简单的编程模型来分布式地处理大数据,其核心是HDFS分布式文件系统和MapReduce分布式计算模型,除此之外,Hadoop生态系统还包括了Hive、HBase、Pig、...

    【企业开源系列】Twitter:收发一条推文的背后

    Apache Pig则作为在Hadoop上的数据分析平台,通过Pig Latin语言简化了对海量数据的MapReduce操作。 其次,Twitter的服务器和存储层依赖于Linux操作系统,为服务器提供了稳定的基础。Memcached作为缓存系统,提升了...

    大数据导论大数据导论大数据导论

    分布式存储计算架构(强烈推荐:Hadoop)分布式程序设计(包含:Apache Pig 或者 Hive)分布式文件系统(比如:Google GFS)多种存储模型,主要包含文档,图,键值,时间序列这几种存储模型(比如:BigTable,Apollo...

    Hadoop和Kerberos:超越大门的疯狂Hadoop and Kerberos: The Madness Beyond the Gate

    它基于Google的MapReduce论文和Google File System (GFS) 论文而设计,最初由Doug Cutting创建,并在2006年作为Apache Lucene的子项目启动。Hadoop的核心组件包括Hadoop Distributed File System (HDFS) 和MapReduce...

    2-大数据技术之Hadoop(入门)

    此外,了解Hadoop与其他大数据工具如Spark、Flink的集成也是很重要的。在实际应用中,根据业务需求选择合适的Hadoop发行版,并掌握集群部署、监控和故障排查技能,能够帮助你有效地利用Hadoop处理和分析大数据。 总...

    IBM数据分析平台方案

    IBM在BigInsights中添加了自家开发的文本分析引擎和数据挖掘工具,以增强对商业分析的支持,并与企业现有系统(如数据库、数据仓库)无缝集成。 5. **BigInsights的版本** 提供基础版和企业版,都包含Apache ...

    hadoop+HBase教程

    首先,Hadoop是一个由Apache软件基金会支持的开源分布式存储与计算框架,其发展起源于Apache Lucene、Apache Nutch以及Google的三大论文:MapReduce、GFS和BigTable。Hadoop生态系统包括Hadoop核心、Hadoop Common、...

    基于Hadoop分布式系统的地质环境大数据框架探讨.pdf

    ArcGIS和MapGIS这样的工具,可以与Hadoop生态系统的其他部分集成,以支持地理空间分析。同时,通过Avro这样的数据序列化系统,可以实现数据高效地在Hadoop组件之间进行传输。 7. 云计算技术 文档提到的云计算技术...

    【最新推荐】hadoop,开题报告-优秀word范文 (8页).pdf

    【Hadoop 开题报告概述】 Hadoop 是一个开源框架,主要设计用于处理和存储海量数据。它最初由 Doug Cutting 领导...此外,Hadoop 的安全性、性能优化以及与其他大数据技术(如Spark)的融合也是未来研究和发展的重点。

    CDHHDPMAPRDKH星环组件比较.pdf

    此外,文档还提到了Apache Falcon、Knox、Phoenix、Pig、Ranger、Slider、Tez、Drill和MapR的特定组件,这些组件分别涉及数据处理、安全管理、SQL支持、YARN应用管理和计算优化等多个方面,构成了丰富的大数据生态...

    SOA通用架构.docx

    这种架构允许不同系统间的组件协同工作,促进业务流程的集成和优化。 在SOA通用架构中,我们通常会涉及到以下几个关键组成部分: 1. **操作系统**:支持SOA架构的操作系统包括Windows、Linux和Unix。这些系统提供...

    大数据技术分享 大数据技术深入浅出 共39页.pdf

    - **概述**:Apache Spark是一个开源集群计算框架,它能够提供比Hadoop MapReduce更快的速度,同时支持多种类型的计算模式。 - **特点**: - **速度**:得益于内存计算和优化的执行计划。 - **通用性**:支持...

Global site tag (gtag.js) - Google Analytics