package testMaven.testMaven;

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer08;
import org.apache.flink.streaming.util.serialization.SimpleStringSchema;

import java.util.Properties;

public class App 
{
    public static void main( String[] args ) throws Exception
    {
        String topic = "test";
        Properties properties = new Properties();
        properties.setProperty("bootstrap.servers", "192.168.1.213:9092");
        properties.setProperty("zookeeper.connect", "192.168.1.213:2181");
        properties.setProperty("group.id", "site");

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.getConfig().enableSysoutLogging();

        DataStream<String> dataStream = env.addSource(new FlinkKafkaConsumer08<String>(topic, new SimpleStringSchema(), properties));

        dataStream.writeAsText("/tmp/app.txt");

        env.execute("read from kafka example");

    }

}
Logo

Kafka开源项目指南提供详尽教程,助开发者掌握其架构、配置和使用,实现高效数据流管理和实时处理。它高性能、可扩展,适合日志收集和实时数据处理,通过持久化保障数据安全,是企业大数据生态系统的核心。

更多推荐