【问题标题】:How to parse json data which comes from kafka topic in Storm scheme class?如何在 Storm 方案类中解析来自 kafka 主题的 json 数据?
【发布时间】:2016-04-04 14:12:54
【问题描述】:

我正在从 kafka 主题中获取 json 数据。 我如何应用 json 解析来获取使用反序列化方法的风暴方案类中所有对象的所有字段,之后我将值返回到新的 return Values().(backtype.storm.tuple.Values 类方法) ?ie,如果我的主题中有 2 个 json 对象,我循环它们以获取所有字段,最后我必须将所有值返回到 return 方法。我的返回应该包含两个 json 对象的所有字段。

我的问题: return 方法中只返回 2 个 obj json 数据。 我认为第二个对象的所有字段都覆盖了第一个对象字段。最后返回第二个对象字段。

你们中的任何人都可以给我一个返回所有对象字段(1,2 个对象字段)的想法......

提前致谢

public class MainParserSpout implements Scheme{
  String tweet_created_at;
  String tweet_id;
  String tweet_id_str;
  String tweet_text;
  String tweet_source;`    
@Override

try{

public List<Object> deserialize(byte[] bytes){
  String twitterEvent = new String(bytes, "UTF-8");
   JSONArray JSON = new JSONArray(twitterEvent);
      for(int i=0;i<JSON.length();i++) {
        JSONObject object_tweet=JSON.getJSONObject(i);
//Tweet status                  
          try{
            this.tweet_created_at=object_tweet.getString("created_at");
            this.tweet_id=object_tweet.getString("id");
            this.tweet_id_str=object_tweet.getString("id_str");
            this.tweet_text=object_tweet.getString("text");
            this.tweet_source=object_tweet.getString("source");
          }catch(Exception e){}
    } //array for close
}catch(Exception e){}
} //JSON array close
  return new Values(tweet_created_at,tweet_id,tweet_id_str,tweet_text,tweet_source);
} //deserialize method close
public Fields getOutputFields() {
    return newFields("tweet_created_at","tweet_id","tweet_id_str","tweet_text","tweet_source");
} //getOutputFields method close
} //class close

【问题讨论】:

  • 我不确定你想做什么...你能举一个小例子来展示你想要得到的两个 JSON 对象和预期的输出元组吗?
  • 我添加了代码示例。如果我的推文对象包含两条推文,则只有第二条推文字段,即:tweet_created_at,id.text,最后返回第二条推文的来源。请分享一个想法如何返回每次迭代的值@Matthias J. Sax
  • 您的代码示例似乎不完整...此外,deserialize 必须返回单个元组。因此,必须将 JSON 中的所有数据收集到单个返回值中。你不能从一条推文中返回多个元组。
  • 是的!我们不能尝试获取多个元组吗?我想要获取多个元组值的方法,可以在storm包中命名任何其他可以帮助我解决这个问题的类。我正在从kafak主题中读取数据。所以我使用了反序列化方法。@Matthias J. Sax
  • 不确定“多个元组值”是什么意思——一个元组有多个值...使用deserialize 是正确的方法;但是,您不能在一次调用中获得多个元组。但是,您可以通过“加倍”您的元组来发出两条推文,即每个值/字段/属性两次。之后,你可以使用一个 Bolt,它接受一个“双推文”,拆分这个元组并发出两个单推文元组。

标签: java apache-kafka apache-storm kafka-consumer-api kafka-producer-api


【解决方案1】:

您不能在一次调用 deserialize 中获得多个元组。但是,您可以通过“加倍”您的元组来发出两条推文,即每个值/字段/属性两次。之后,你可以使用一个 Bolt,它接受一个“双推文”,拆分这个元组并发出两个单推文元组。

类似的东西(我不熟悉 JSON Tweet 格式,所以这是关于问题的代码示例的更多猜测):

@Override
public List<Object> deserialize(byte[] bytes){
  List<String> doubleTweet = new ArrayList<String>();

  try{
    String twitterEvent = new String(bytes, "UTF-8");
    JSONArray JSON = new JSONArray(twitterEvent);


    for(int i=0;i<JSON.length();i++) {
      JSONObject object_tweet=JSON.getJSONObject(i);
      for(int j=0;j<object_tweet.length();j++){
        //Tweet status                  
        try{
          doubleTweet.add(object_tweet.getString("created_at"));
          doubleTweet.add(object_tweet.getString("id"));
          doubleTweet.add(object_tweet.getString("id_str"));
          doubleTweet.add(object_tweet.getString("text"));
          doubleTweet.add(object_tweet.getString("source"));
        }catch(Exception e){}
      }
    }
  }catch(Exception e){}

  return doubleTweet;
}

doubleTweet 包含每个字段两次(第一个推文的字段 0-4 和第二个推文的字段 5-9)。因此,一个连续的螺栓可以只提取这些字段并为每条推文发出一个 5 字段元组)。

作为替代方案,您也可以使用RawScheme 并在后续螺栓中进行 JSON 解析。在这个螺栓中,您可以发出多个元组(即,每条推文一个):https://github.com/apache/storm/tree/master/external/storm-kafka#multischeme

如果您使用 RawScheme,则 spout 会发出带有单个 byte[] 字段的元组。因此,您可以在 Bolt.execute() 中执行 JSON 解析,并为每条推文调用 Collector.emit()。

【讨论】:

  • 你能给我一个“加倍”和“RawScheme”的代码示例吗?@Matthias J. Sax
  • 是否需要返回列表对象“doubleTweet”才能返回newFields()?对问题做了一个小的更正,请看一下@Matthias J. Sax
  • 当我将列表对象返回给新的 feilds();我收到以下错误,如下所示 Error:java.lang.IllegalArgumentException: Tuple created with wrong number of fields。预期 1 个字段,但得到 5 个字段。@Matthias J. Sax
  • 我从不单独使用 Kafka 和 Schema。不确定,为什么它需要一个字段而不是五个字段。您可以将 Stacktrace 添加到您的问题中吗?以及组装拓扑的代码?
【解决方案2】:

我错过了 kafka 是消息发布-订阅消息系统这一点。 当我尝试将数据发送给生产者时,我将 Json 卡盘 20 个对象作为单个消息发送,但我的方案仅适用于单个 Json 卡盘。所以我将单个 20 个对象 Json 卡盘分成 20 个 json 卡盘并发送每个转给 Json 制作人。

【讨论】:

    猜你喜欢
    • 2016-11-16
    • 2021-07-15
    • 2016-11-23
    • 2021-03-20
    • 1970-01-01
    • 2016-05-27
    • 1970-01-01
    • 2017-11-17
    • 2016-08-27
    相关资源
    最近更新 更多