使用Flume采集Twitter数据至Kafka时遇配置及API权限问题求助
解决Flume采集Twitter数据到Kafka的两类问题
一、修复Source/Sink与MemChannel的配置不兼容错误
错误原因
首次启动报错核心原因:
- Flume官方Twitter Source默认
batchSize为1000,Kafka Sink默认batchSize也为1000 - 你的MemChannel配置中
transactionCapacity仅为100,小于Source和Sink的batch size,违反Flume约束(Source/Sink的batch size不能超过Channel的transactionCapacity)
解决方法
两种方案任选其一:
方案1:提升MemChannel的transactionCapacity
将transactionCapacity调整为不小于Source和Sink的batch size(比如1000),同时保证capacity(总容量)大于transactionCapacity(你的配置中capacity=10000已满足要求)。修改后的Channel配置片段:
TwitterAgent.channels.MemChannel.type = memory TwitterAgent.channels.MemChannel.capacity = 10000 TwitterAgent.channels.MemChannel.transactionCapacity = 1000
方案2:降低Source和Sink的batchSize
显式设置Twitter Source和Kafka Sink的batchSize为不超过100(原transactionCapacity值),示例:
# 给Twitter Source添加batchSize配置 TwitterAgent.sources.Twitter.batchSize = 100 # Kafka Sink已设置batchSize=1,无需修改 TwitterAgent.sinks.kafkasink.batchSize = 1
二、修复Twitter API 403权限错误
错误原因
你使用的官方Flume Twitter Source基于twitter4j库,该库调用的是Twitter API v1.1流式接口,但你的开发者账号仅拥有Essential权限,这类权限只能访问Twitter API v2端点,因此被拒绝访问(错误码453)。
解决方法
方案1:申请Elevated权限
在Twitter开发者门户提交Elevated权限申请,填写使用场景(比如社交媒体数据采集与分析),审核通过后即可获得API v1.1访问权限,满足官方Twitter Source的需求。权限生效后,确保Flume配置中的密钥正确,重启Flume即可。
方案2:自定义支持API v2的Flume Source
如果不想申请Elevated权限,可以基于Twitter API v2的流式接口开发自定义Flume Source。API v2的Essential权限允许访问部分流式端点(比如过滤流),可参考Twitter API v2文档,使用官方v2 SDK或HTTP请求实现数据采集,封装为Flume Source插件。
修改后的完整配置示例(采用方案1调整transactionCapacity)
# Licensed to the Apache Software Foundation (ASF) under one # or more contributor license agreements. See the NOTICE file # distributed with this work for additional information # regarding copyright ownership. The ASF licenses this file # to you under the Apache License, Version 2.0 (the # "License"); you may not use this file except in compliance # with the License. You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, # software distributed under the License is distributed on an # "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. # The configuration file needs to define the sources, # the channels and the sinks. # Sources, channels and sinks are defined per agent, # in this case called 'TwitterAgent' TwitterAgent.sources = Twitter TwitterAgent.channels = MemChannel TwitterAgent.sinks = kafkasink #filesink TwitterAgent.sources.Twitter.type = org.apache.flume.source.twitter.TwitterSource TwitterAgent.sources.Twitter.channels = MemChannel TwitterAgent.sources.Twitter.keywords = AWS Marketplace, GCP Marketplace, Azure Marketplace #kafkasink TwitterAgent.sinks.kafkasink.type = org.apache.flume.sink.kafka.KafkaSink TwitterAgent.sinks.kafkasink.topic = fp_mc TwitterAgent.sinks.kafkasink.brokerList = fp_mc-kafka-1:9092 TwitterAgent.sinks.kafkasink.channel = MemChannel TwitterAgent.sinks.kafkasink.batchSize = 1 TwitterAgent.channels.MemChannel.type = memory TwitterAgent.channels.MemChannel.capacity = 10000 TwitterAgent.channels.MemChannel.transactionCapacity = 1000 TwitterAgent.sources.Twitter.consumerKey = xxxx TwitterAgent.sources.Twitter.consumerSecret = xxxx TwitterAgent.sources.Twitter.accessToken = xxxx TwitterAgent.sources.Twitter.accessTokenSecret = xxxx
内容的提问来源于stack exchange,提问作者netdevmike
相关产品推荐
相关产品推荐

