ラベル java の投稿を表示しています。 すべての投稿を表示
ラベル java の投稿を表示しています。 すべての投稿を表示

2013年12月7日土曜日

Google Guava's EventBus programming example.

GoogleGuava には EventBus を簡単に実装できるライブラリが含まれていたので、試してみたので備忘録的に残します。

Publish/Subscribe を使ってEDP(Event Driven Programming)を簡単に実装できるようです。

それではさっそくシンプルに

pub/sub するときに受け渡しするメッセージオブジェクト(POJO)
public class EventMessage {

    final int msgcode;
    final String msg;

    public EventMessage(int msgcode, String msg) {
 this.msgcode = msgcode;
 this.msg = msg;
    }

    public int getMsgcode() {
 return msgcode;
    }

    public String getMsg() {
 return msg;
    }
}

そして、subscriber、@Subscribeアノテーションが付けられたメソッドが呼ばれるようです。
なんとなく、1クラスで1つのSubscribe定義がよさそうな感じ...

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.google.common.eventbus.Subscribe;

public class Subscriber {
 
 private static Logger log = LoggerFactory.getLogger(Subscriber.class);

 protected static ExecutorService worker = Executors.newCachedThreadPool();
    
 public static class EventHandler implements Runnable {
 
  final EventMessage message;
 
  public EventHandler( EventMessage message) {
   this.message = message;
  }
 
  public void run() {
   try {
    log.debug("[start] Processing worker thread. / {} {}", message.msgcode, message.msg);
    TimeUnit.MILLISECONDS.sleep(message.msgcode); // do something
    log.debug("[end] worker thread");
   } catch (Exception e) {
    log.error("Event failed", e);
   }
  }
    }
    
 @Subscribe
 public void handleEvent(final EventMessage eventMessage) {
  try {
   log.debug("[start] dispatched event message to worker thread. {}/{}", eventMessage.msgcode, eventMessage.msg);
   worker.execute(new EventHandler(eventMessage));
   log.debug("[end] dispatched events.");
  } catch (Exception e) {
   log.error("event failed.", e);
  }
 }
    
 public void shutdown() {
  worker.shutdown();
  try {
   if (!worker.awaitTermination(1000, TimeUnit.MILLISECONDS))
    worker.shutdownNow();
   log.info("shutdown success");
  } catch (InterruptedException e) {
   log.error("shutdown failed",e);
  }
 }
}

最後に、Main というか Publish するところ、前にあるSubscriberの受け取るメッセージオブジェクトに対応しないメッセージがとどくとDeadEventsSubscriberがハンドルする。

import java.util.concurrent.TimeUnit;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.google.common.eventbus.DeadEvent;
import com.google.common.eventbus.EventBus;
import com.google.common.eventbus.Subscribe;

public class Main {
 private static Logger log = LoggerFactory.getLogger(Main.class);

 public static class DeadEventsSubscriver {
  @Subscribe
  public void handleDeadEvent(DeadEvent deadEvent) {
   log.error("DEAD EVENT: {}", deadEvent.getEvent());
  }
 }

 /**
 * example EventBus programming 'Google-Guava'
 * @param args
 */
 public static void main(String[] args) {
 try {
  final EventBus eventBus = new EventBus("example-events");
  final Subscriber subscriver = new Subscriber();
  final DeadEventsSubscriver des = new DeadEventsSubscriver();
      
  eventBus.register(subscriver);
  eventBus.register(des);

  // published
  eventBus.post("This is dead event");
  eventBus.post(new EventMessage(5000, "This message from hoge"));
  eventBus.post(new EventMessage(1000, "This message from fuga"));

  log.debug("waiting...");
  TimeUnit.MILLISECONDS.sleep(10000);

  log.info(">> end of main thread");
  eventBus.unregister(subscriver);
  eventBus.unregister(des);
  subscriver.shutdown();
      
 } catch (InterruptedException e) {
  log.error("fail", e);
 }
 }
}

んで、実行した結果。
[main] 20:44:28.803  ERROR o.horiga.study.eventbus.example.Main[20] - DEAD EVENT: This is dead event
[main] 20:44:28.807  DEBUG o.h.s.eventbus.example.Subscriber[40] - [start] dispatched event message to worker thread. 5000/This message from hoge
[main] 20:44:28.808  DEBUG o.h.s.eventbus.example.Subscriber[42] - [end] dispatched events.
[main] 20:44:28.809  DEBUG o.h.s.eventbus.example.Subscriber[40] - [start] dispatched event message to worker thread. 1000/This message from fuga
[pool-1-thread-1] 20:44:28.809  DEBUG o.h.s.eventbus.example.Subscriber[28] - [start] Processing worker thread. / 5000 This message from hoge
[main] 20:44:28.809  DEBUG o.h.s.eventbus.example.Subscriber[42] - [end] dispatched events.
[pool-1-thread-2] 20:44:28.810  DEBUG o.h.s.eventbus.example.Subscriber[28] - [start] Processing worker thread. / 1000 This message from fuga
[main] 20:44:28.810  DEBUG o.horiga.study.eventbus.example.Main[42] - waiting...
[pool-1-thread-2] 20:44:29.811  DEBUG o.h.s.eventbus.example.Subscriber[30] - [end] worker thread
[pool-1-thread-1] 20:44:33.810  DEBUG o.h.s.eventbus.example.Subscriber[30] - [end] worker thread
[main] 20:44:38.811  INFO  o.horiga.study.eventbus.example.Main[45] - &gt&gt end of main thread
[main] 20:44:38.816  INFO  o.h.s.eventbus.example.Subscriber[54] - shutdown success


すげー簡単にでけた。ここまでホントに15分程度。

とりあえず、eventbus#post は、blocking されるようだったので、subscriber側は thread にしてみたこれで asynchronous event driven っぽくなったかな? subscribe をnetwork経由で別のsubscriberに通知できればもっといい感じになりそうだけどそのあたりは、akka とか使ったほうが早いだろうかと

まだ、スレッドセーフとか渡したいEventのPojoとかいろいろな型でためしてどれくらい拡張性(POJO となるEventクラスの継承をした場合、@Subscribeが複数ある場合は対象となる?とか、AnnotationやらGenericsなどを使ったりとか)があるか試したいところですが、ここまでということで。

今回試したコードはこちら

2013年6月14日金曜日

【備忘録】spring + aspectj によるAOP

AOP (Aspect Oriented Programming) はある特定の振る舞い(aspect)を分離し、既存の振る舞いに対し入れ込むような時に利用される。例えば以下のようなものが挙げられます。

  • 特定処理にログを入れる
  • 開始終了の処理時間を計測する
  • 特定処理の呼び出し回数を計測する
今回は、spring mvc を使ってある特定の controller の処理結果をハンドリングし、例外が発生した場合、自動的にエラーオブジェクトを生成して結果を返すような処理を作ってみたので、備忘録的にまとめておきます。

controller は以下、入力を受けて service の処理結果をjsonとして返却します。
@Controller
public class HelloController {
  
  private static Logger logger = LoggerFactory.getLogger(HelloController.class);
  
  @Autowired
  private HelloService service;
  
  @Procedure
  @RequestMapping(value="/hello/{type}", method={RequestMethod.GET})
  public @ResponseBody
  Output<?> hello(
      @PathVariable("type") String type) throws Exception {
    logger.info("- start hello");
    try {
    if ("exception".equalsIgnoreCase(type))
      return new Output<Hello>(UUID.randomUUID().toString(), service.runOnException());
    else if ("failure".equalsIgnoreCase(type)) 
      return new Output<Hello>(UUID.randomUUID().toString(), service.runOnFailure());
    
    return new Output<Hello>(UUID.randomUUID().toString(), service.runOnSuccess());
    } finally {
      logger.info("- end hello");
    }
  }
}

output は以下のようなpojoです。
public class Output<T> {
  
  private String trxId;
  private int statusCode;
  private String statusMessage;
  private T data;

  public Output( String trxId, int statusCode, String statusMessage) {
    this.trxId = trxId;
    this.statusCode = statusCode;
    this.statusMessage = "";
  }
  
  public Output( String trxId, T data) {
    this.trxId = trxId;
    this.statusCode = 200;
    this.data = data;
    this.statusMessage = "";
  }
  
  public String getTrxId() {
    return trxId;
  }

  public void setTrxId(String trxId) {
    this.trxId = trxId;
  }

  public int getStatusCode() {
    return statusCode;
  }

  public void setStatusCode(int statusCode) {
    this.statusCode = statusCode;
  }

  public String getStatusMessage() {
    return statusMessage;
  }

  public void setStatusMessage(String statusMessage) {
    this.statusMessage = statusMessage;
  }

  public T getData() {
    return data;
  }

  public void setData(T data) {
    this.data = data;
  }
  
  public String toString() {
    return new StringBuilder().append("@").append(this.trxId).append("-[")
        .append(this.statusCode).append(": ")
        .append(this.statusMessage).append("], ").append(data)
        .toString();
  }
}

で、AOP を利用しない場合、呼び出し結果は以下のようになります。
10:52:40.811 [http-8080-2] INFO  j.b.h.example.aop.HelloController - - start hello
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.HelloServiceImpl - success
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.HelloController - - end hello

このままでは例外発生時にエラー情報を返却できないです。
なので、hello() に対し、@Procedure というアノテーションを付けて、@Procedure がついた処理について例外が発生した場合、エラーのOutputを生成しレスポンスしたいと考えてみました。
AOP は AspectJ を利用します。Spring+AspectJ はこちらを参考に
@Component
@Aspect
public class ProcedureAspect {

  static Logger logger = LoggerFactory.getLogger(ProcedureAspect.class);

  @Pointcut("execution(* jp.blogspot.horiga3.*.*(..)) ")
  public void targetMethods() {}
  
  @Before("@annotation(jp.blogspot.horiga3.example.aop.Procedure)")
  public void preHandle() {
    logger.info("Aspect :: preHandle");
  }
  
  @AfterReturning(
      pointcut="@annotation(jp.blogspot.horiga3.example.aop.Procedure)",
      returning="retVal")
  public void postHandle(Object retVal) {
    logger.info("Aspect :: postHandle, retVal={}", retVal != null ? retVal.toString() : "null");
  }
  
  @Around("@annotation(jp.blogspot.horiga3.example.aop.Procedure)")
  public Object handle(ProceedingJoinPoint pjp) {

    logger.info("Aspect :: around - start");

    Object[] args;
    try {
      args = pjp.getArgs();
      return args == null ? pjp.proceed() : pjp.proceed(args);
    } catch (Throwable e) {
      logger.info("Aspect :: handleException");
      int statusCode = 500;
      String statusMessage = "unknown";
      if (e instanceof ProcedureException) {
        statusCode = ((ProcedureException) e).getStatusCode();
        statusMessage = ((ProcedureException) e).getStatusMessage();
      } else if (e instanceof IllegalArgumentException) {
        statusCode = 400;
        statusMessage = "Invalid parameter";
      }
      Output<Object> error = new Output<Object>(UUID.randomUUID().toString(), statusCode, statusMessage);
      return error;
    } finally {
      logger.info("Aspect :: around - end");
    }
  }
}

spring の設定は、aop:aspectj-autoproxy を追加しただけ
<?xml version="1.0" encoding="UTF-8"?>
<beans 
  xmlns="http://www.springframework.org/schema/beans"
  xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xmlns:aop="http://www.springframework.org/schema/aop" 
  xmlns:mvc="http://www.springframework.org/schema/mvc"
  xmlns:context="http://www.springframework.org/schema/context"
  xsi:schemaLocation=
     "http://www.springframework.org/schema/mvc http://www.springframework.org/schema/mvc/spring-mvc.xsd
    http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
    http://www.springframework.org/schema/aop http://www.springframework.org/schema/aop/spring-aop.xsd
    http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context.xsd" >
  
  <mvc:annotation-driven />

  <aop:aspectj-autoproxy/>
  <context:component-scan base-package="jp.blogspot.horiga3.example.aop"/>
  
</beans>

で、実行すると例外の場合もjsonを返却するようになりました。既存のControllerは何も修正していないですね。
10:52:40.768 [http-8080-2] INFO  j.b.h.example.aop.ProcedureAspect - Aspect :: around - start
10:52:40.772 [http-8080-2] INFO  j.b.h.example.aop.ProcedureAspect - Aspect :: preHandle
10:52:40.811 [http-8080-2] INFO  j.b.h.example.aop.HelloController - - start hello
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.HelloServiceImpl - success
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.HelloController - - end hello
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.ProcedureAspect - Aspect :: around - end
10:52:40.845 [http-8080-2] INFO  j.b.h.example.aop.ProcedureAspect - Aspect :: postHandle, retVal=@11d8c89a-6fc2-4e38-a745-f2ade9c3d6ff-[200: ], jp.blogspot.horiga3.example.aop.Hello@6dbf4a72

@AfterReturningが @Around より後に来ることは予想できたけど、@Around が @Before より先に来るんですね。

2013年6月8日土曜日

springframework を利用した JavaMail 送信 の覚え書き

springframework を利用して Mail を送信することがあったので、覚え書き。spring は何かと設定を xml にする必要があって覚えるのが大変と思って Guice をここ2年程度利用していたが、最近は annotation でほぼ設定できるようになってて結構良い感じだった。まぁ、springframework はそれだけではなく巨大なフレームワークだからいろいろと機能が沢山あって全てを語るには勉強が足りないですww。

個人的に最近気になっているのは、playframework と、vertx かなと JVM + netty をベースにした framework が良い性能をだしていて少しづつ勉強しているところです。

話がずれたので、とりあえず覚え書き。

まず、spring の applicationContext.xml で設定する bean 設定。mail 設定の部分だけ別ファイルとして設定
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
 xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
 xsi:schemaLocation="http://www.springframework.org/schema/beans 
 http://www.springframework.org/schema/beans/spring-beans.xsd">
 
 <bean id="mailSender" 
  class="org.springframework.mail.javamail.JavaMailSenderImpl">
  <property name="host" value="${smtp.host}"/>
 </bean>
 
 <bean id="velocityEngine" class="org.springframework.ui.velocity.VelocityEngineFactoryBean">
  <property name="resourceLoaderPath" value="classpath:mail" />
  <property name="velocityPropertiesMap">
   <map>
                <entry key="input.encoding" value="UTF-8" />
                <entry key="output.encoding" value="UTF-8" />
            </map>
  </property>
 </bean>
 
 <bean id="velocityJavaMailSender" class="jp.blogspot.horiga3.example.spring.mail.VelocityJavaMailSender">
  <property name="mailSender" ref="mailSender" />
  <property name="velocityEngine" ref="velocityEngine" />
 </bean>
</beans>

SMTPサーバはないと動きません。幸いにも社内にあるのでそれを設定します。それぞれの環境に合わせて変更するひつようがありますのであしからず。

あとは、自前クラスは以下のように
import javax.mail.internet.MimeMessage;

import org.apache.velocity.app.VelocityEngine;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.mail.javamail.JavaMailSender;
import org.springframework.mail.javamail.MimeMessageHelper;
import org.springframework.mail.javamail.MimeMessagePreparator;
import org.springframework.stereotype.Component;
import org.springframework.ui.velocity.VelocityEngineUtils;

@Component
public class VelocityJavaMailSender {

 @Autowired
 protected JavaMailSender mailSender;
 
 @Autowired
 protected VelocityEngine velocityEngine;

 public static class MailMessage {
  
  private boolean html = false;
  
  private String from;
  private String personal;
  private String mailTemplate;
  private String[] recipients;
  private String subject;
  private Map<String, Object> content;

  public boolean isHtml() {
   return html;
  }

  public void setHtml(boolean html) {
   this.html = html;
  }

  public String getFrom() {
   return from;
  }

  public void setFrom(String from) {
   this.from = from;
  }

  public String getPersonal() {
   return personal;
  }

  public void setPersonal(String personal) {
   this.personal = personal;
  }

  public String getMailTemplate() {
   return mailTemplate;
  }

  public void setMailTemplate(String mailTemplate) {
   this.mailTemplate = mailTemplate;
  }

  public String[] getRecipients() {
   return recipients;
  }

  public void setRecipients(String[] recipients) {
   this.recipients = recipients;
  }

  public String getSubject() {
   return subject;
  }

  public void setSubject(String subject) {
   this.subject = subject;
  }

  public Map<String, Object> getContent() {
   return content;
  }

  public void setContent(Map<String, Object> content) {
   this.content = content;
  }
 }
 
 public void sendMailMessage(MailMessage mail) throws Exception {
  this.mailSender.send(createMailPreparator(mail));
 }
 
 public void setMailSender(JavaMailSender mailSender) {
  this.mailSender = mailSender;
 }
 
 public void setVelocityEngine(VelocityEngine velocityEngine) {
  this.velocityEngine = velocityEngine;
 }
 
 private MimeMessagePreparator createMailPreparator(
   final MailMessage mailMessage) {
  MimeMessagePreparator preparator = new MimeMessagePreparator() {
   @Override
   public void prepare(MimeMessage mimeMessage) throws Exception {
    MimeMessageHelper message = new MimeMessageHelper(mimeMessage);
    message.setTo(mailMessage.getRecipients());
    message.setSubject(mailMessage.getSubject());
    if ( null != mailMessage.getPersonal() && mailMessage.getPersonal().trim().length() > 0)
     message.setFrom(mailMessage.getFrom(), mailMessage.getPersonal());
    else message.setFrom(mailMessage.getFrom());
    message.setText(VelocityEngineUtils.mergeTemplateIntoString(
      velocityEngine, mailMessage.getMailTemplate(), "utf-8",
      mailMessage.getContent()), mailMessage.isHtml());
   }
  };
  return preparator;
 }
}


で、テストケースが以下。
import java.util.HashMap;

import junit.framework.Assert;

import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

import jp.blogspot.horiga3.example.spring.mail.VelocityJavaMailSender.MailMessage;


@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration(locations={
 "file:src/main/webapp/WEB-INF/springframework/mail-context.xml"
})
public class VelocityJavaMailSenderTest {
 
 @Autowired
 @Qualifier("velocityJavaMailSender")
 VelocityJavaMailSender mailSender;
 
 @Test
 public void test() {
  
  MailMessage message = new MailMessage();
  
  try {
   message.setMailTemplate("test.vm");
   message.setRecipients(new String[] { "your mail address " });
   message.setSubject("test message");
   message.setFrom("noreply@hogehoge.com");
   message.setPersonal("JUnitさん");
   HashMap<String, Object> content = new HashMap<String, Object>();
   content.put("str", "test");
   content.put("n", System.currentTimeMillis()/1000);
   message.setContent(content);
   mailSender.sendMailMessage(message);
  } catch ( Exception e) {
   e.printStackTrace();
   Assert.fail(e.getMessage());
  }
 }
}

送信先を自分のメールアドレスに設定してテストケースを実行すると
正しく送信されていることが確認できた。
ちなみに少しハマったところが、テンプレートファイルのエンコーディングとかも全て統一しておかないと文字化けします。
テストケースにこのソースがあるとmavenビルドとかで、自動テストするとビルドのたびに毎回メールくるようになります。注意しましょう〜







2013年5月14日火曜日

Consistent Hashing による 分散テスト


大量データを分散させる技術?アルゴリズムで consistent hashing ということが挙げられます。mixi さんのエンジニアブログでも紹介されていましたが、Javaでサンプルコードを書いてみました。

package com.blogspot.horiga3.example;

import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.SortedMap;
import java.util.TreeMap;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;

import com.google.common.hash.HashFunction;
import com.google.common.hash.Hashing;

public class ConsistentHashing<T> {

 static final HashFunction DEFAULT_HASH_FUNCTION = Hashing.md5();
 static final int DEFAULT_NUMBER_OF_REPLICAS = 200;
 
 static int numberOfReplicas = DEFAULT_NUMBER_OF_REPLICAS;
 
 final SortedMap<Long, T> shardedNodes = new TreeMap<Long, T>();
 final HashFunction hashFunction;
 
 public ConsistentHashing(Collection<T> nodes, HashFunction hashFunction) {
  this.hashFunction = hashFunction != null ? hashFunction : DEFAULT_HASH_FUNCTION;
  for (T node : nodes) {
   for (int i = 0; i < numberOfReplicas; i++) {
    shardedNodes.put(hashFunction.hashString(String.format("SHARDS-%s-%03d", node.toString(), i)).asLong(), node);
   }
  }
 }

 public T get(Object key) {
  if (shardedNodes.isEmpty()) {
   return null;
  }
  long hash = hashFunction.hashString(key.toString()).asLong();
  if (!shardedNodes.containsKey(hash)) {
   SortedMap<Long, T> tailMap = shardedNodes.tailMap(hash);
   hash = tailMap.isEmpty() ? shardedNodes.firstKey() : tailMap.firstKey();
  }
  return shardedNodes.get(hash);
 }

 public static void main(String[] args) {
  try {
   List<String> shards = new ArrayList<String>();
   
   final int shardsCount = 5;
   final int numberOfUsers = 1000000;
   
   for (int i = 1; i <= shardsCount; i++)
    shards.add("shard" + i);

   List<String> keys = new ArrayList<String>();
   for (int i = 0; i < numberOfUsers; i++) {
    keys.add(UUID.randomUUID().toString().replaceAll("-", ""));
   }
   
   final HashFunction hashFunction = Hashing.md5();
   ConsistentHashing<String> consistentHashing = new ConsistentHashing<String>(shards, hashFunction);

   SortedMap<String, List<String>> keyshards = new TreeMap<String, List<String>>();
   for (String key : keys) {
    String shard = consistentHashing.get(key);
    if (keyshards.containsKey(shard))
     keyshards.get(shard).add(key);
    else {
     keyshards.put(shard, new ArrayList<String>());
     keyshards.get(shard).add(key);
    }
   }
   
   // shard node added
   shards.add(String.format("shard%d", shardsCount + 1));

   ConsistentHashing<String> consistentHashing2 = new ConsistentHashing<String>(shards, hashFunction);
   SortedMap<String, List<String>> keyshards2 = new TreeMap<String, List<String>>();
   for (String key : keys) {
    String k = consistentHashing2.get(key);
    if (keyshards2.containsKey(k))
     keyshards2.get(k).add(key);
    else {
     keyshards2.put(k, new ArrayList<String>());
     keyshards2.get(k).add(key);
    }
   }

   SortedMap<String, AtomicInteger> result = new TreeMap<String, AtomicInteger>();
   final String unchanged = "shard.unchanged";
   for (String key : keys) {
    String shard1 = consistentHashing.get(key);
    String shard2 = consistentHashing2.get(key);
    if (shard1.equals(shard2)) {
     if (!result.containsKey(unchanged))
      result.put(unchanged, new AtomicInteger(1));
     else
      result.get(unchanged).incrementAndGet();
    } else {
     String k = shard1 + "=>" + shard2;
     if (!result.containsKey(k))
      result.put(k, new AtomicInteger(1));
     else
      result.get(k).incrementAndGet();
    }
   }
   
   System.out.println("==========================================");
   System.out.println(":: Consistent hashing sharding Testcase ::");
   
   System.out.println("------------------------------------------");
   for (Map.Entry<String, List<String>> entry: keyshards.entrySet()) {
    System.out.println(String.format("%s: %d", entry.getKey(), entry.getValue().size()));
   }
   System.out.println("------------------------------------------ :: number of key for after adding a shard");
   for (Map.Entry<String, List<String>> entry: keyshards2.entrySet()) {
    if ( !entry.getKey().equals(String.format("shard%d", shardsCount + 1))) {
     int before = keyshards.get(entry.getKey()).size();
     int after = entry.getValue().size();
     System.out.println(String.format("%s: %d=>%d(-%d)", entry.getKey(), before, after, before-after));
    } else {
     System.out.println(String.format("%s: 0=>%d", entry.getKey(), entry.getValue().size()));
    }
   }
   System.out.println("------------------------------------------ :: number of key shard node that has moved");
   for (Map.Entry<String, AtomicInteger> entry : result.entrySet()) {
    if (!entry.getKey().equals(unchanged)) {
     System.out.println(String.format("%s: %d", entry.getKey(), entry.getValue().intValue()));
    }
   }

   System.out.println("=====================================");
   System.out.println("shard.unchanged=" + result.get(unchanged).intValue());
   System.out.println("shard.changeing=" + (keys.size() - result.get(unchanged).intValue()));
   System.out.println("-------------------------------------");
   System.out.println("HIT's=" + Math.round(result.get(unchanged).intValue() * 100 / keys.size()) + "%");

  } catch (Exception e) {
   e.printStackTrace();
  }
 }
}


テストは、ノードが5台から1台追加した場合の、100万件のキーの移動についてテストしてみたけど、ヒット率は、およそ80%程度、キーの移動はおよそ20%程度になりました。
==========================================
:: Consistent hashing sharding Testcase ::
------------------------------------------
shard1: 185436
shard2: 218328
shard3: 207460
shard4: 170488
shard5: 218288
------------------------------------------ :: number of key for after adding a shard
shard1: 185436=>144072(-41364)
shard2: 218328=>185908(-32420)
shard3: 207460=>177256(-30204)
shard4: 170488=>144561(-25927)
shard5: 218288=>179190(-39098)
shard6: 0=>169013
------------------------------------------ :: number of key shard node that has moved
shard1=>shard6: 41364
shard2=>shard6: 32420
shard3=>shard6: 30204
shard4=>shard6: 25927
shard5=>shard6: 39098
=====================================
shard.unchanged=830987
shard.changeing=169013
-------------------------------------
HIT's=83%

また、上記サンプルコードを利用して数パターン試した結果は以下のようになりました
キーの数変更前node数変更後node数key移動件数ヒット率
10000005616901383%
10000003434443365%
1000056169283%
1000034244875%

Consistent Hashing を利用することで、ノードを追加しても、既存のノード間でのキーの移動は考慮しなくてよさそうなので、良い方法でした。

2013年5月3日金曜日

Template Engine - mustache : How to import another template.

以前にメモした {{mustache}} テンプレートについて書きましたが、mustache を使って対象のテンプレートに共通のヘッダやらフッターやらをimportしたいって思いますよね。
以下のようにできます。まぁ本家のマニュアルにも記載してありましたが,,

まずは、元になるテンプレートを用意します。
import したいテンプレートの部分は {{ > 対象となるテンプレートのパス }} のように設定します。

{{> common/common_header }}
  {{!This is comment of Mustache!!}}
  <h1>pojo.str</h1>
  <h2>{{str}}</h2>
  
  <h1>pojo.num</h1>
  <h2>{{num}}</h2>
  
  <h1>pojo.flag</h1>
  {{#flag}}flag = true {{/flag}}
  {{^flag}}flag = false {{/flag}}
  
  <h1>pojo.array</h1>
  {{#array}}
  {{.}}</br>
  {{/array}}
  
  <h1>pojo.data</h1>
  {{#data}}
  <p>upper: {{a}} , {{b}}</p></br>
  {{/data}}
  
  <h1>Escaped Characters</h1>
  {{escape}}
{{> common/common_footer }}

で import されるテンプレートは以下のように定義

<!DOCTYPE html>
<html>
<head>
<title>{{title}}</title>
<meta charset="UTF-8">
</head><body>

<p>This is footer</p>
</body>
</html>

プログラム側は、以下のように ※前回のとほぼ変わりません。テンプレートに渡すオブジェクトを2つ渡しているだけです

package com.blogspot.agiroh.netty.template.html;

import java.io.PrintWriter;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;

import com.github.mustachejava.DefaultMustacheFactory;
import com.github.mustachejava.Mustache;
import com.github.mustachejava.MustacheFactory;

public class ExampleMustache {

 final static String template_path = "template/mustache/example.mustache";

 public static class ExamplePojo {
  String str;
  long num;
  boolean flag;
  Map<String, Object> data;
  List<String> array;
  String escape;
 }

 static Map<String, Mustache> templates = new HashMap<>();

 public void run() throws Exception {

  // default is classpath.
  MustacheFactory factory = new DefaultMustacheFactory(); 
  
  Mustache mustache = null;
  if (!templates.containsKey(template_path)) {
   mustache = factory.compile(template_path);
   templates.put(template_path, mustache);
  } else {
   System.out.println("templates from cache");
   mustache = templates.get(template_path);
  }

  ExamplePojo pojo = new ExamplePojo();
  pojo.str = UUID.randomUUID().toString();
  pojo.num = new Date().getTime() / 1000;
  pojo.flag = true;
  pojo.data = new HashMap<String, Object>();
  pojo.data.put("a", "A");
  pojo.data.put("b", "B");
  pojo.array = new ArrayList<>();
  pojo.array.add("hoge");
  pojo.array.add("fuga");
  pojo.escape = "<p>\"te&st\"</p>";

  Map<String, String> headerPojo = new HashMap<>();
  headerPojo.put("title", "Import another templtes");

  mustache.execute(new PrintWriter(System.out), new Object[] { pojo, headerPojo }).flush();
 }

 public static void main(String[] args) {
  try {
   new ExampleMustache().run();
  } catch (Exception e) {
   e.printStackTrace();
  }
 }
}

これで思ったとおり動作できました。

2013年3月10日日曜日

Netty with html template engine.


最近 playframework やら、twitter の Finagle などが利用されていることで知られる nettyをいろいろ個人的に使っているのですが、netty は基本的に network framework なので基本的には html を生成するような昨日はありません。こういったことはアプリケーションの実装に任されています。playframework などは基盤のevent-drivenのnetworkフレームワークにnettyを利用しながら、httpのframeworkとしてよく考えられていることで人気があります。
playframework を使っても良かったのですが、playframework も設定などはplayframeworkの作法があります。netty だけを使ってアプリケーションは自由に実装したいと思いまずは、html をレンダリングすることからとおもいます。html のテンプレートエンジンは前回のブログで書いた mustache を使おうとおもいます。

まずは、pom.xml
<project 
  xmlns="http://maven.apache.org/POM/4.0.0" 
  xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
 <modelVersion>4.0.0</modelVersion>

 <groupId>com.blogspot.3agiroh.netty.template.html</groupId>
 <artifactId>mastache-template-test</artifactId>
 <version>0.0.1-SNAPSHOT</version>
 <packaging>jar</packaging>

 <name>mastache-template-test</name>
 <url>http://maven.apache.org</url>

 <properties>
  <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
 </properties>

 <build>
  <plugins>
   <plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-compiler-plugin</artifactId>
    <configuration>
     <source>1.7</source>
     <target>1.7</target>
     <encoding>UTF-8</encoding>
    </configuration>
   </plugin>
   <plugin>
    <groupId>org.apache.felix</groupId>
    <artifactId>maven-bundle-plugin</artifactId>
    <version>2.3.7</version>
    <extensions>true</extensions>
    <configuration>
     <instructions>
      <!--<Export-Package>*</Export-Package> <Import-Package>*</Import-Package> -->
     </instructions>
    </configuration>
   </plugin>
  </plugins>
 </build>

 <dependencies>
  <!-- mustache template engine. 
  @see https://github.com/spullara/mustache.java -->
  <dependency>
   <groupId>com.github.spullara.mustache.java</groupId>
   <artifactId>compiler</artifactId>
   <version>0.8.10</version>
  </dependency>

  <dependency>
   <groupId>io.netty</groupId>
   <artifactId>netty</artifactId>
   <version>3.6.2.Final</version>
  </dependency>
  <dependency>
   <groupId>org.slf4j</groupId>
   <artifactId>slf4j-api</artifactId>
   <version>1.7.2</version>
  </dependency>
  <dependency>
   <groupId>ch.qos.logback</groupId>
   <artifactId>logback-classic</artifactId>
   <version>1.0.9</version>
  </dependency>
  <dependency>
   <groupId>commons-configuration</groupId>
   <artifactId>commons-configuration</artifactId>
   <version>1.9</version>
  </dependency>
  <dependency>
   <groupId>commons-lang</groupId>
   <artifactId>commons-lang</artifactId>
   <version>2.6</version>
  </dependency>

  <dependency>
   <groupId>junit</groupId>
   <artifactId>junit</artifactId>
   <version>3.8.1</version>
   <scope>test</scope>
  </dependency>
 </dependencies>
</project>

起動のクラスはこんな感じで

package com.blogspot.agiroh.netty.template.html;

import java.net.InetSocketAddress;
import java.util.concurrent.Executors;

import org.jboss.netty.bootstrap.ServerBootstrap;
import org.jboss.netty.channel.ChannelPipeline;
import org.jboss.netty.channel.ChannelPipelineFactory;
import org.jboss.netty.channel.Channels;
import org.jboss.netty.channel.socket.nio.NioServerSocketChannelFactory;
import org.jboss.netty.handler.codec.http.HttpChunkAggregator;
import org.jboss.netty.handler.codec.http.HttpContentCompressor;
import org.jboss.netty.handler.codec.http.HttpRequestDecoder;
import org.jboss.netty.handler.codec.http.HttpResponseEncoder;

public class SimpleHttpServer {

 final int port;
 
 public SimpleHttpServer( int port) {
  this.port = port;
 }
 
 public void run() {
  
  ServerBootstrap bootstrap = new ServerBootstrap(
      new NioServerSocketChannelFactory(
        Executors.newCachedThreadPool(), 
        Executors.newCachedThreadPool()));
  
  bootstrap.setPipelineFactory(new ChannelPipelineFactory() {
   
   @Override
   public ChannelPipeline getPipeline() throws Exception {
    
    ChannelPipeline pipeline = Channels.pipeline();
    pipeline.addLast("decoder", new HttpRequestDecoder());
    pipeline.addLast("aggregator", new HttpChunkAggregator(1048576));
    pipeline.addLast("encoder", new HttpResponseEncoder());
    pipeline.addLast("deflater", new HttpContentCompressor());
    pipeline.addLast("htmlrender", new HtmlRenderHandler());
    
    return pipeline;
   }
  });
  
  bootstrap.bind(new InetSocketAddress(port));
  
 }
 
 public static void main(String[] args) {
  int port = 9000;
  new SimpleHttpServer(port).run();
 }
}
で、HTMLを生成するところは mustache java を利用してこんな感じ
package com.blogspot.agiroh.netty.template.html;

import java.io.IOException;
import java.io.PrintWriter;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

import org.apache.commons.lang.StringUtils;
import org.jboss.netty.buffer.ChannelBuffer;
import org.jboss.netty.buffer.ChannelBufferOutputStream;
import org.jboss.netty.buffer.ChannelBuffers;
import org.jboss.netty.channel.ChannelFuture;
import org.jboss.netty.channel.ChannelFutureListener;
import org.jboss.netty.channel.ChannelHandlerContext;
import org.jboss.netty.channel.ExceptionEvent;
import org.jboss.netty.channel.MessageEvent;
import org.jboss.netty.channel.SimpleChannelUpstreamHandler;
import org.jboss.netty.handler.codec.http.DefaultHttpResponse;
import org.jboss.netty.handler.codec.http.HttpHeaders;
import org.jboss.netty.handler.codec.http.HttpMessage;
import org.jboss.netty.handler.codec.http.HttpRequest;
import org.jboss.netty.handler.codec.http.HttpResponse;
import org.jboss.netty.handler.codec.http.HttpResponseStatus;
import org.jboss.netty.handler.codec.http.HttpVersion;
import org.jboss.netty.handler.codec.http.QueryStringDecoder;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.github.mustachejava.DefaultMustacheFactory;
import com.github.mustachejava.Mustache;
import com.github.mustachejava.MustacheFactory;

public class HtmlRenderHandler
  extends SimpleChannelUpstreamHandler {

 static Logger log = LoggerFactory.getLogger(HtmlRenderHandler.class);

 public static final int RES_BUFFER_CAPACITY = 65535;

 public static class ExamplePojo {
  String str;
  long num;
  boolean flag;
  Map<String, Object> data;
  List<String> array;
  String escape;
 }

 @Override
 public void messageReceived(
   ChannelHandlerContext ctx, MessageEvent e) throws Exception {
  
  Object msg = e.getMessage();
  
  if (msg instanceof HttpRequest) {

   HttpRequest request = (HttpRequest) msg;

   if (HttpHeaders.is100ContinueExpected(request)) {
    send100Continue(e);
   }

   ExamplePojo pojo = new ExamplePojo();
   pojo.str = "テスト";
   pojo.num = 30;
   pojo.flag = true;

   QueryStringDecoder qsd = new QueryStringDecoder(request.getUri());
   Map<String, List<String>> params = qsd.getParameters();
   pojo.data = new HashMap<String, Object>();
   if (!params.isEmpty()) {
    for (Map.Entry<String, List<String>> entry : params.entrySet())
     pojo.data.put(entry.getKey(), StringUtils.join(entry.getValue(), ","));
   } else {
    pojo.data.put("key1", "value1");
   }

   pojo.array = Arrays.asList("hoge", "fuga");
   pojo.escape = "<p>hogehogehoge</p>";

   render(e, "html/ja/test.mustache", pojo);
   return;
  }
  
  log.debug("[unknown]");
 }

 protected void render(
   MessageEvent e,  // HTTP request event.
   String resource, // mustache template. 
   Object data)   // mustache template bindings.
     throws IOException {
  
  log.debug("[{}] start rendering.", resource);
  
  HttpMessage req = (HttpMessage) e.getMessage();
  HttpResponse res = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK);

  MustacheFactory mf = new DefaultMustacheFactory();
  Mustache mustache = mf.compile(resource);

  ChannelBuffer buffer = ChannelBuffers.buffer(RES_BUFFER_CAPACITY);
  ChannelBufferOutputStream out = new ChannelBufferOutputStream(buffer);
  mustache.execute(new PrintWriter(out), data).flush();
  res.setContent(out.buffer());
  try {
   out.close();
  } catch (Exception ee) {}

  res.setHeader(HttpHeaders.Names.CONTENT_TYPE, "text/html; charset=utf-8;");

  final boolean keepalive = HttpHeaders.isKeepAlive(req);
  if (keepalive) {
   res.setHeader(HttpHeaders.Names.CONTENT_LENGTH, res.getContent().readableBytes());
   res.setHeader(HttpHeaders.Names.CONNECTION, HttpHeaders.Values.KEEP_ALIVE);
  }

  // Unsupported cookie.

  ChannelFuture future = e.getChannel().write(res);
  if (!keepalive) {
   future.addListener(ChannelFutureListener.CLOSE);
  }

 }

 private static void send100Continue(MessageEvent e) {
  log.info("send 100 continue.");
  HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.CONTINUE);
  e.getChannel().write(response);
 }
 
 @Override
    public void exceptionCaught(ChannelHandlerContext ctx, ExceptionEvent e)
            throws Exception {
  log.error("Error!! : {}", e.getCause().getMessage());
        e.getCause().printStackTrace();
        e.getChannel().close();
    }
}

そしてテンプレートは以下のようにして実行する。

<!DOCTYPE html>
<html>
<head>
<title>Mustache</title>
<meta charset="UTF-8">
</head>
<body>
 {{!This is comment of Mustache!!}}
 <h1>pojo.str</h1>
 <h2>{{str}}</h2>
 
 <h1>pojo.num</h1>
 <h2>{{num}}</h2>
 
 <h1>pojo.flag</h1>
 {{#flag}}flag = true {{/flag}}
 {{^flag}}flag = false {{/flag}}
 
 <h1>pojo.array</h1>
 {{#array}}
 {{.}}</br>
 {{/array}}
 
 <h1>pojo.data</h1>
 {{#data}}
 <p>upper: {{a}} , {{b}}</p></br>
 {{/data}}
 
 <h1>Escaped Characters</h1>
 {{escape}}
</body>
</html>

http://localhost:9000?a=hogehoge&b=fugafuga こんなアクセスをすると結果こうなりました。きちんと表示されますね。


Template Engine - mustache

テンプレートエンジンはみなさんご存知かと思います。有名なのですと以下のようなものがあります。
  • JSP
  • Velocity
  • FreeMarker
「テンプレートエンジンはテンプレートと呼ばれる雛形と、あるデータモデルで表現される入力データを合成し、成果ドキュメントを出力するソフトウェアまたはソフトウェアコンポーネントである」とあります。

最近だと以下の様なのあるようです。

  • Thymeleaf
    • 所見:コーダーさんが用意してくれるhtmlをほぼ変更せずにテンプレートエンジンを適用できるような感じ。
  • Mustache
    • 所見:「ますたっしゅ」って呼ぶそうです。Ruby、Java、Python、Javascript、node.js、... と多くの言語をサポート。使い方がシンプル。
今回は mustache についてJavaを使って簡単に試してみた結果をまとめます。
まずはテンプレートは以下の様なものを用意しました。

<!DOCTYPE html>
<html>
<head>
<title>Mustache</title>
<meta charset="UTF-8">
</head>
<body>
 {{!This is comment of Mustache!!}}
 <h1>pojo.str</h1>
 <h2>{{str}}</h2>
 
 <h1>pojo.num</h1>
 <h2>{{num}}</h2>
 
 <h1>pojo.flag</h1>
 {{#flag}}flag = true {{/flag}}
 {{^flag}}flag = false {{/flag}}
 
 <h1>pojo.array</h1>
 {{#array}}
 {{.}}</br>
 {{/array}}
 
 <h1>pojo.data</h1>
 {{#data}}
 <p>upper: {{a}} , {{b}}</p></br>
 {{/data}}
 
 <h1>Escaped Characters</h1>
 {{escape}}
</body>
</html>
</html>

{{ }} で囲われたところにデータで埋め込まれます。{{#xxx}}〜{{/xxx}}が配列データの繰り返しです。囲われたデータにアクセスする際は、{{.}}とドットでアクセスします。またMapデータのようなkey-valueは、{{key}}でvalue要素にアクセスします。


でJavaは以下のように用意してみました。

package com.blogspot.agiroh.netty.template.html;

import java.io.PrintWriter;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;

import com.github.mustachejava.DefaultMustacheFactory;
import com.github.mustachejava.Mustache;
import com.github.mustachejava.MustacheFactory;

public class ExampleMustache {
 
 final static String template_path = "template/mustache/example.mustache";
 
 public static class ExamplePojo {
  String str;
  long num;
  boolean flag;
  Map<String, Object> data;
  List<String> array;
  String escape;
 }
 
 static Map<String, Mustache> templates = new HashMap<>();
 
 public void run() throws Exception {
  
  MustacheFactory factory = new DefaultMustacheFactory(); // default is classpath root
  
  Mustache mustache = null;
  if ( !templates.containsKey(template_path)) {
   
   mustache = factory.compile(template_path);
   templates.put(template_path, mustache);
  } else {
   System.out.println("templates from cache");
   mustache = templates.get(template_path);
  }
  
  ExamplePojo pojo = new ExamplePojo();
  pojo.str = UUID.randomUUID().toString();
  pojo.num = new Date().getTime()/1000;
  pojo.flag = true;
  pojo.data = new HashMap<String, Object>();
  pojo.data.put("a", "A");
  pojo.data.put("b", "B");
  pojo.array = new ArrayList<>();
  pojo.array.add("hoge");
  pojo.array.add("fuga");
  pojo.escape = "<p>\"te&st\"</p>";
  
  mustache.execute( new PrintWriter(System.out), pojo).flush();
 }
 
 public static void main(String[] args) {
  try {
    new ExampleMustache().run();
  } catch (Exception e) {
   e.printStackTrace();
  }
 } 
}

結果は以下になります。きちんとhtmlのエスケープ処理もされています。なかなか使いやすいとおもいますがいかがでしょうか。

<!DOCTYPE html>
<html>
<head>
<title>Mustache</title>
<meta charset="UTF-8">
</head>
<body>
 
 <h1>pojo.str</h1>
 <h2>5b954508-cf30-4501-b173-75a33446e4ba</h2>
 
 <h1>pojo.num</h1>
 <h2>1362906356</h2>
 
 <h1>pojo.flag</h1>
 flag = true 
 
 
 <h1>pojo.array</h1>
 hoge</br>
 fuga</br>
 
 <h1>pojo.data</h1>
 <p>upper: A , B</p></br>
 
 <h1>Escaped Characters</h1>
 &lt;p&gt;&quot;te&amp;st&quot;&lt;/p&gt;
</body>
</html>
</html>

2013年3月8日金曜日

fluent-logger-java を使ってみた

アプリケーションログを fluentd で収集して、アプリケーションエラーのログとかイベントログの収集してKPIとか集めたくてちょっとまとめます。
とりあえず、ローカル環境は MacOS ですが、fluentd のインストールについては割愛します。rubyやらなんやら入れる必要あります。
※ここら辺を参照

まずは、fluentd 側の設定。本家サイトの通りでおk


[horiga@fluent]: cat fluent4j.conf 
<source></source>
 type forward
 port 24224

<match fluentd.test.**>
 type stdout
</match>


で、起動コマンドはこんな感じで


[horiga@bin]: ./fluentd --config ../fluent4j.conf 


問題なく起動して、コンソールにログが出力されます

次はアプリケーション側。


<project xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xmlns="http://maven.apache.org/POM/4.0.0"
  xsi:schemalocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
 <modelversion>4.0.0</modelversion>

 <groupid>com.blogspot.3agiroh.examples.fluentlogger</groupid>
 <artifactid>fluentlogger</artifactid>
 <version>0.0.1-SNAPSHOT</version>
 <packaging>jar</packaging>

 <name>fluentlogger</name>
 <url>http://maven.apache.org</url>

 <properties>
  <project .build.sourceencoding="">UTF-8</project>
 </properties>

 <repositories>
  <repository>
   <id>fluentd.org</id>
   <name>Fluentd Maven2 Repository</name>
   <url>http://fluentd.org/maven2</url>
  </repository>
 </repositories>

 <dependencies>

  <dependency>
      <groupid>com.google.code.gson</groupid>
      <artifactid>gson</artifactid>
      <version>2.2.2</version>
      <scope>compile</scope>
    </dependency>
    
  <dependency>
   <groupid>org.fluentd</groupid>
   <artifactid>fluent-logger</artifactid>
   <version>0.2.6</version>
   
  </dependency>

  <dependency>
   <groupid>junit</groupid>
   <artifactid>junit</artifactid>
   <version>3.8.1</version>
   <scope>test</scope>
  </dependency>
 </dependencies>
</project>


今回テストしたのはJavaでこんなコードで実行しました
本家サイトにありますが、少しバージョンが古かったですね。Map に文字列だけでなく数字、bool型も


package com.blogspot.agiroh.examples.fluentlogger;

import java.util.Date;
import java.util.HashMap;
import java.util.Map;

import org.fluentd.logger.FluentLogger;

import com.google.gson.Gson;

/**
 * This is Fluentd4 java test 
 */
public class App {

 static FluentLogger logger = FluentLogger.getLogger("fluentd.test");
 
 public static class Pojo {
  String str;
  int num;
  
  public String toString() {
   return new Gson().toJson(this);
  }
 }
 
 public static void main(String[] args) {
   Map<string object=""> data = new HashMap<string object="">();
   data.put("name", "hoge");
   data.put("age", 35);
   data.put("loggedin", true);
   long ts = new Date().getTime();
   data.put("ts", ts);
   
   Pojo pojo = new Pojo();
   pojo.str = "fuga";
   pojo.num = 100;
   data.put("pojo", pojo);
   
   // default timestamp is [sec]
   logger.log("fluent4j", data);
   
   // timestamp [millis]
   // logger.log("fluent4j", data, new Date().getTime());
 }
}


結果


[horiga@bin]: ./fluentd --config ../fluent4j.conf 
2013-03-08 20:05:11 +0900: starting fluentd-0.10.23
2013-03-08 20:05:11 +0900: reading config file path="../fluent4j.conf"
2013-03-08 20:05:11 +0900: adding source type="forward"
2013-03-08 20:05:11 +0900: adding match pattern="fluentd.test.**" type="stdout"
2013-03-08 20:05:11 +0900: listening fluent socket on 0.0.0.0:24224
2013-03-08 20:13:59 +0900 fluentd.test.fluent4j: {"ts":1362741239520,"loggedin":true,"age":35,"name":"hoge","pojo":"{\"str\":\"fuga\",\"num\":100}"}


おお!!無事に fluentd さんにイベントが転送されました。※本来ならこれを収集サーバに forward すればいいと思う。

でも、やっぱりvalueにpojoはだめか。まぁ複雑にしないでkey-valueの方が使いやすいからこれは問題ないかな。これで accesslog 以外にもいろいろ使えることはおk。実際のロジックと集計などは別に分離したいよね。